diff --git a/core/src/main/scala/kafka/server/BrokerServer.scala b/core/src/main/scala/kafka/server/BrokerServer.scala index 76023d5153870..3507b88c94028 100644 --- a/core/src/main/scala/kafka/server/BrokerServer.scala +++ b/core/src/main/scala/kafka/server/BrokerServer.scala @@ -312,23 +312,6 @@ class BrokerServer( config, Some(clientToControllerChannelManager), None, None, groupCoordinator, transactionCoordinator) - /* Add all reconfigurables for config change notification before starting the metadata listener */ - config.dynamicConfig.addReconfigurables(this) - - dynamicConfigHandlers = Map[String, ConfigHandler]( - ConfigType.Topic -> new TopicConfigHandler(logManager, config, quotaManagers, None), - ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers)) - - if (!config.processRoles.contains(ControllerRole)) { - // If no controller is defined, we rely on the broker to generate snapshots. - metadataSnapshotter = Some(new BrokerMetadataSnapshotter( - config.nodeId, - time, - threadNamePrefix, - new BrokerSnapshotWriterBuilder(raftManager.client) - )) - } - metadataListener = new BrokerMetadataListener( config.nodeId, time, @@ -369,9 +352,6 @@ class BrokerServer( featuresRemapped ) - // Register a listener with the Raft layer to receive metadata event notifications - raftManager.register(metadataListener) - val endpoints = new util.ArrayList[Endpoint](networkListeners.size()) var interBrokerListener: Endpoint = null networkListeners.iterator().forEachRemaining(listener => { @@ -408,6 +388,26 @@ class BrokerServer( }.toMap } + /* Add all reconfigurables for config change notification before starting the metadata listener */ + config.dynamicConfig.addReconfigurables(this) + + dynamicConfigHandlers = Map[String, ConfigHandler]( + ConfigType.Topic -> new TopicConfigHandler(logManager, config, quotaManagers, None), + ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers)) + + if (!config.processRoles.contains(ControllerRole)) { + // If no controller is defined, we rely on the broker to generate snapshots. + metadataSnapshotter = Some(new BrokerMetadataSnapshotter( + config.nodeId, + time, + threadNamePrefix, + new BrokerSnapshotWriterBuilder(raftManager.client) + )) + } + + // Register a listener with the Raft layer to receive metadata event notifications + raftManager.register(metadataListener) + val fetchManager = new FetchManager(Time.SYSTEM, new FetchSessionCache(config.maxIncrementalFetchSessionCacheSlots, KafkaServer.MIN_INCREMENTAL_FETCH_SESSION_EVICTION_MS)) diff --git a/core/src/main/scala/kafka/server/ControllerServer.scala b/core/src/main/scala/kafka/server/ControllerServer.scala index 4f2ad837b8351..b823dca5489d4 100644 --- a/core/src/main/scala/kafka/server/ControllerServer.scala +++ b/core/src/main/scala/kafka/server/ControllerServer.scala @@ -24,8 +24,8 @@ import kafka.network.{DataPlaneAcceptor, SocketServer} import kafka.raft.KafkaRaftManager import kafka.security.CredentialProvider import kafka.server.KafkaConfig.{AlterConfigPolicyClassNameProp, CreateTopicPolicyClassNameProp} -import kafka.server.KafkaRaftServer.BrokerRole import kafka.server.QuotaFactory.QuotaManagers +import kafka.server.metadata.DynamicConfigPublisher import kafka.utils.{CoreUtils, Logging} import kafka.zk.{KafkaZkClient, ZkMigrationClient} import org.apache.kafka.common.config.ConfigException @@ -35,6 +35,7 @@ import org.apache.kafka.common.security.token.delegation.internals.DelegationTok import org.apache.kafka.common.utils.LogContext import org.apache.kafka.common.{ClusterResource, Endpoint} import org.apache.kafka.controller.{Controller, QuorumController, QuorumFeatures} +import org.apache.kafka.image.publisher.MetadataPublisher import org.apache.kafka.metadata.KafkaConfigSchema import org.apache.kafka.metadata.authorizer.ClusterMetadataAuthorizer import org.apache.kafka.metadata.bootstrap.BootstrapMetadata @@ -104,6 +105,7 @@ class ControllerServer( var quotaManagers: QuotaManagers = _ var controllerApis: ControllerApis = _ var controllerApisHandlerPool: KafkaRequestHandlerPool = _ + def kafkaYammerMetrics: KafkaYammerMetrics = KafkaYammerMetrics.INSTANCE var migrationSupport: Option[ControllerMigrationSupport] = None private def maybeChangeStatus(from: ProcessStatus, to: ProcessStatus): Boolean = { @@ -118,13 +120,6 @@ class ControllerServer( true } - private def doRemoteKraftSetup(): Unit = { - // Explicitly configure metric reporters on this remote controller. - // We do not yet support dynamic reconfiguration on remote controllers in general; - // remove this once that is implemented. - new DynamicMetricReporterState(config.nodeId, config, metrics, clusterId) - } - def clusterId: String = sharedServer.metaProps.clusterId def startup(): Unit = { @@ -242,11 +237,17 @@ class ControllerServer( case _ => // nothing to do } controller = controllerBuilder.build() + val metadataPublishers = new java.util.ArrayList[MetadataPublisher]() - // Perform any setup that is done only when this node is a controller-only node. - if (!config.processRoles.contains(BrokerRole)) { - doRemoteKraftSetup() - } + val dynamicConfigHandlers = Map[String, ConfigHandler]( + ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers) + ) + metadataPublishers.add(new DynamicConfigPublisher( + config, + sharedServer.metadataPublishingFaultHandler, + dynamicConfigHandlers.toMap, + "controller" + )) if (config.migrationEnabled) { val zkClient = KafkaZkClient.createZkClient("KRaft Migration", time, config, KafkaServer.zkClientConfigFromKafkaConfig(config)) @@ -290,6 +291,14 @@ class ControllerServer( s"${DataPlaneAcceptor.MetricPrefix}RequestHandlerAvgIdlePercent", DataPlaneAcceptor.ThreadPrefix) + config.dynamicConfig.addReconfigurables(this) + + // Install the metadata publishers. Note that they will not actually receive any metadata + // until we catch up to the high water mark of __cluster_metadata-0. We are not waiting for + // that. We are just waiting for the installation process to complete. + FutureUtils.waitWithLogging(logger.underlying, "all of the metadata publishers to be installed", + sharedServer.loader.installPublishers(metadataPublishers), startupDeadline, time) + /** * Enable the controller endpoint(s). If we are using an authorizer which stores * ACLs in the metadata log, such as StandardAuthorizer, we will be able to start diff --git a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala index b924648c691cd..af69228deda2d 100755 --- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala +++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala @@ -271,6 +271,35 @@ class DynamicBrokerConfig(private val kafkaConfig: KafkaConfig) extends Logging addBrokerReconfigurable(new DynamicProducerStateManagerConfig(kafkaServer.logManager.producerStateManagerConfig)) } + /** + * Add reconfigurables to be notified when a dynamic controller config is updated. + */ + def addReconfigurables(controller: ControllerServer): Unit = { + controller.authorizer match { + case Some(authz: Reconfigurable) => addReconfigurable(authz) + case _ => + } + if (kafkaConfig.isKRaftCoResidentMode) { + debug("Relying on the broker to reconfigure the Yammer metrics.") + } else { + addReconfigurable(controller.kafkaYammerMetrics) + } + if (kafkaConfig.isKRaftCoResidentMode) { + debug("Relying on the broker to reconfigure the metric reporters.") + } else { + addReconfigurable(new DynamicMetricsReporters( + kafkaConfig.brokerId, controller.config, controller.metrics, controller.clusterId)) + } + addBrokerReconfigurable(controller.socketServer) + + // KAFKA-14349: add dynamic thread pool resizing here + + // KAFKA-14350: add dynamic listener reconfiguration here + + // When we implement controller mutation quotas in KRaft, as discussed in KAFKA-14351, + // we'll need to make them reconfigurable here (probably via DynamicClientQuotaCallback). + } + def addReconfigurable(reconfigurable: Reconfigurable): Unit = CoreUtils.inWriteLock(lock) { verifyReconfigurableConfigs(reconfigurable.reconfigurableConfigs.asScala) reconfigurables.add(reconfigurable) diff --git a/core/src/main/scala/kafka/server/KafkaRaftServer.scala b/core/src/main/scala/kafka/server/KafkaRaftServer.scala index d11fbe99aa6d9..d31ee6db522b3 100644 --- a/core/src/main/scala/kafka/server/KafkaRaftServer.scala +++ b/core/src/main/scala/kafka/server/KafkaRaftServer.scala @@ -77,10 +77,7 @@ class KafkaRaftServer( ) private val broker: Option[BrokerServer] = if (config.processRoles.contains(BrokerRole)) { - Some(new BrokerServer( - sharedServer, - offlineDirs - )) + Some(new BrokerServer(sharedServer, offlineDirs)) } else { None } diff --git a/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala b/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala index 43d79e88601b4..0e37215682800 100644 --- a/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala +++ b/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala @@ -211,7 +211,7 @@ class BrokerMetadataPublisher( } // Apply configuration deltas. - dynamicConfigPublisher.publish(delta, newImage) + dynamicConfigPublisher.publish(delta, newImage, null) // Apply client quotas delta. try { diff --git a/core/src/main/scala/kafka/server/metadata/DynamicConfigPublisher.scala b/core/src/main/scala/kafka/server/metadata/DynamicConfigPublisher.scala index 12ff51d4039f3..95cce37a32798 100644 --- a/core/src/main/scala/kafka/server/metadata/DynamicConfigPublisher.scala +++ b/core/src/main/scala/kafka/server/metadata/DynamicConfigPublisher.scala @@ -22,6 +22,7 @@ import kafka.server.ConfigAdminManager.toLoggableProps import kafka.server.{ConfigEntityName, ConfigHandler, ConfigType, KafkaConfig} import kafka.utils.Logging import org.apache.kafka.common.config.ConfigResource.Type.{BROKER, TOPIC} +import org.apache.kafka.image.loader.LoaderManifest import org.apache.kafka.image.{MetadataDelta, MetadataImage} import org.apache.kafka.server.fault.FaultHandler @@ -31,10 +32,16 @@ class DynamicConfigPublisher( faultHandler: FaultHandler, dynamicConfigHandlers: Map[String, ConfigHandler], nodeType: String -) extends Logging { +) extends org.apache.kafka.image.publisher.MetadataPublisher with Logging { logIdent = s"[DynamicConfigPublisher nodeType=${nodeType} id=${conf.nodeId}] " - def publish(delta: MetadataDelta, newImage: MetadataImage): Unit = { + def name(): String = "DynamicConfigPublisher" + + def publish( + delta: MetadataDelta, + newImage: MetadataImage, + manifest: LoaderManifest + ): Unit = { val deltaName = s"MetadataDelta up to ${newImage.highestOffsetAndEpoch().offset}" try { // Apply configuration deltas. diff --git a/core/src/test/scala/integration/kafka/server/KRaftClusterTest.scala b/core/src/test/scala/integration/kafka/server/KRaftClusterTest.scala index d554dc7b37777..bfd6c31cb2d24 100644 --- a/core/src/test/scala/integration/kafka/server/KRaftClusterTest.scala +++ b/core/src/test/scala/integration/kafka/server/KRaftClusterTest.scala @@ -24,6 +24,7 @@ import kafka.utils.TestUtils import org.apache.kafka.clients.admin.AlterConfigOp.OpType import org.apache.kafka.clients.admin._ import org.apache.kafka.common.acl.{AclBinding, AclBindingFilter} +import org.apache.kafka.common.config.ConfigException import org.apache.kafka.common.config.ConfigResource import org.apache.kafka.common.config.ConfigResource.Type import org.apache.kafka.common.message.DescribeClusterRequestData @@ -31,7 +32,7 @@ import org.apache.kafka.common.network.ListenerName import org.apache.kafka.common.protocol.Errors._ import org.apache.kafka.common.quota.{ClientQuotaAlteration, ClientQuotaEntity, ClientQuotaFilter, ClientQuotaFilterComponent} import org.apache.kafka.common.requests.{ApiError, DescribeClusterRequest, DescribeClusterResponse} -import org.apache.kafka.common.{Endpoint, TopicPartition, TopicPartitionInfo} +import org.apache.kafka.common.{Endpoint, Reconfigurable, TopicPartition, TopicPartitionInfo} import org.apache.kafka.image.ClusterImage import org.apache.kafka.metadata.BrokerState import org.apache.kafka.server.authorizer._ @@ -39,12 +40,15 @@ import org.apache.kafka.server.common.MetadataVersion import org.apache.kafka.server.log.remote.storage.RemoteLogManagerConfig import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{Tag, Test, Timeout} +import org.junit.jupiter.params.ParameterizedTest +import org.junit.jupiter.params.provider.ValueSource import org.slf4j.LoggerFactory import java.io.File import java.nio.file.{FileSystems, Path} +import java.util.concurrent.atomic.AtomicInteger import java.{lang, util} -import java.util.concurrent.CompletionStage +import java.util.concurrent.{CompletableFuture, CompletionStage} import java.util.{Arrays, Collections, Optional, OptionalLong, Properties} import scala.annotation.nowarn import scala.collection.mutable @@ -966,6 +970,47 @@ class KRaftClusterTest { cluster.close() } } + + @ParameterizedTest + @ValueSource(booleans = Array(false, true)) + def testReconfigureControllerAuthorizer(combinedMode: Boolean): Unit = { + val cluster = new KafkaClusterTestKit.Builder( + new TestKitNodes.Builder(). + setNumBrokerNodes(1). + setCoResident(combinedMode). + setNumControllerNodes(1).build()). + setConfigProp("authorizer.class.name", classOf[FakeConfigurableAuthorizer].getName). + build() + + def assertFoobarValue(expected: Int): Unit = { + TestUtils.retry(60000) { + assertEquals(expected, cluster.controllers().values().iterator().next(). + authorizer.get.asInstanceOf[FakeConfigurableAuthorizer].foobar.get()) + assertEquals(expected, cluster.brokers().values().iterator().next(). + authorizer.get.asInstanceOf[FakeConfigurableAuthorizer].foobar.get()) + } + } + + try { + cluster.format() + cluster.startup() + cluster.waitForReadyBrokers() + assertFoobarValue(0) + val admin = Admin.create(cluster.clientProperties()) + try { + admin.incrementalAlterConfigs( + Collections.singletonMap(new ConfigResource(Type.BROKER, ""), + Collections.singletonList(new AlterConfigOp( + new ConfigEntry(FakeConfigurableAuthorizer.foobarConfigKey, "123"), OpType.SET)))). + all().get() + } finally { + admin.close() + } + assertFoobarValue(123) + } finally { + cluster.close() + } + } } class BadAuthorizer() extends Authorizer { @@ -986,3 +1031,71 @@ class BadAuthorizer() extends Authorizer { override def deleteAcls(requestContext: AuthorizableRequestContext, aclBindingFilters: util.List[AclBindingFilter]): util.List[_ <: CompletionStage[AclDeleteResult]] = ??? } + +object FakeConfigurableAuthorizer { + val foobarConfigKey = "fake.configurable.authorizer.foobar.config" + + def fakeConfigurableAuthorizerConfigToInt(configs: util.Map[String, _]): Int = { + val result = configs.get(foobarConfigKey) + if (result == null) { + 0 + } else { + val resultString = result.toString().trim() + try { + Integer.valueOf(resultString) + } catch { + case e: NumberFormatException => throw new ConfigException(s"Bad value of ${foobarConfigKey}: ${resultString}") + } + } + } +} + +class FakeConfigurableAuthorizer() extends Authorizer with Reconfigurable { + import FakeConfigurableAuthorizer._ + + val foobar = new AtomicInteger(0) + + override def start(serverInfo: AuthorizerServerInfo): java.util.Map[Endpoint, _ <: CompletionStage[Void]] = { + serverInfo.endpoints().asScala.map(e => e -> { + val future = new CompletableFuture[Void] + future.complete(null) + future + }).toMap.asJava + } + + override def reconfigurableConfigs(): java.util.Set[String] = Set(foobarConfigKey).asJava + + override def validateReconfiguration(configs: util.Map[String, _]): Unit = { + fakeConfigurableAuthorizerConfigToInt(configs) + } + + override def reconfigure(configs: util.Map[String, _]): Unit = { + foobar.set(fakeConfigurableAuthorizerConfigToInt(configs)) + } + + override def authorize(requestContext: AuthorizableRequestContext, actions: util.List[Action]): util.List[AuthorizationResult] = { + actions.asScala.map(_ => AuthorizationResult.ALLOWED).toList.asJava + } + + override def acls(filter: AclBindingFilter): lang.Iterable[AclBinding] = List[AclBinding]().asJava + + override def close(): Unit = {} + + override def configure(configs: util.Map[String, _]): Unit = { + foobar.set(fakeConfigurableAuthorizerConfigToInt(configs)) + } + + override def createAcls( + requestContext: AuthorizableRequestContext, + aclBindings: util.List[AclBinding] + ): util.List[_ <: CompletionStage[AclCreateResult]] = { + Collections.emptyList() + } + + override def deleteAcls( + requestContext: AuthorizableRequestContext, + aclBindingFilters: util.List[AclBindingFilter] + ): util.List[_ <: CompletionStage[AclDeleteResult]] = { + Collections.emptyList() + } +} diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/LoaderManifest.java b/metadata/src/main/java/org/apache/kafka/image/loader/LoaderManifest.java new file mode 100644 index 0000000000000..60889997f3ce4 --- /dev/null +++ b/metadata/src/main/java/org/apache/kafka/image/loader/LoaderManifest.java @@ -0,0 +1,36 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.kafka.image.loader; + +import org.apache.kafka.image.MetadataProvenance; + + +/** + * Contains information about what was loaded. + */ +public interface LoaderManifest { + /** + * Describes the type of manifest which this is. + */ + LoaderManifestType type(); + + /** + * The highest offset and epoch included in the new image, inclusive. + */ + MetadataProvenance provenance(); +} diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/LoaderManifestType.java b/metadata/src/main/java/org/apache/kafka/image/loader/LoaderManifestType.java new file mode 100644 index 0000000000000..f83cadc8a26e7 --- /dev/null +++ b/metadata/src/main/java/org/apache/kafka/image/loader/LoaderManifestType.java @@ -0,0 +1,27 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.kafka.image.loader; + + +/** + * Contains information about the type of a loader manifest. + */ +public enum LoaderManifestType { + LOG_DELTA, + SNAPSHOT; +} diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/LogDeltaManifest.java b/metadata/src/main/java/org/apache/kafka/image/loader/LogDeltaManifest.java index 982a1f8e27180..6a8cec4e4cfc2 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/LogDeltaManifest.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/LogDeltaManifest.java @@ -26,7 +26,7 @@ /** * Contains information about a set of changes that were loaded from the metadata log. */ -public class LogDeltaManifest { +public class LogDeltaManifest implements LoaderManifest { /** * The highest offset and epoch included in this delta, inclusive. */ @@ -66,7 +66,12 @@ public LogDeltaManifest( this.numBytes = numBytes; } + @Override + public LoaderManifestType type() { + return LoaderManifestType.LOG_DELTA; + } + @Override public MetadataProvenance provenance() { return provenance; } diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java index 21df1a761cd5e..50014dda979a8 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/MetadataLoader.java @@ -258,7 +258,8 @@ private void maybeInitializeNewPublishers() { try { log.info("Publishing initial snapshot at offset {} to {}", image.highestOffsetAndEpoch().offset(), publisher.name()); - publisher.publishSnapshot(delta, image, manifest); + publisher.publish(delta, image, manifest); + publisher.handleControllerChange(currentLeaderAndEpoch); publishers.put(publisher.name(), publisher); } catch (Throwable e) { faultHandler.handleFault("Unhandled error publishing the initial metadata " + @@ -295,7 +296,7 @@ public void handleCommit(BatchReader reader) { log.debug("Publishing new image with provenance {}.", image.provenance()); for (MetadataPublisher publisher : publishers.values()) { try { - publisher.publishLogDelta(delta, image, manifest); + publisher.publish(delta, image, manifest); } catch (Throwable e) { faultHandler.handleFault("Unhandled error publishing the new metadata " + "image ending at " + manifest.provenance().lastContainedOffset() + @@ -392,7 +393,7 @@ public void handleSnapshot(SnapshotReader reader) { log.debug("Publishing new snapshot image with provenance {}.", image.provenance()); for (MetadataPublisher publisher : publishers.values()) { try { - publisher.publishSnapshot(delta, image, manifest); + publisher.publish(delta, image, manifest); } catch (Throwable e) { faultHandler.handleFault("Unhandled error publishing the new metadata " + "image from snapshot at offset " + reader.lastContainedLogOffset() + @@ -449,6 +450,15 @@ SnapshotManifest loadSnapshot( public void handleLeaderChange(LeaderAndEpoch leaderAndEpoch) { eventQueue.append(() -> { currentLeaderAndEpoch = leaderAndEpoch; + for (MetadataPublisher publisher : publishers.values()) { + try { + publisher.handleControllerChange(currentLeaderAndEpoch); + } catch (Throwable e) { + faultHandler.handleFault("Unhandled error publishing the new leader " + + "change to " + currentLeaderAndEpoch + " with publisher " + + publisher.name(), e); + } + } }); } diff --git a/metadata/src/main/java/org/apache/kafka/image/loader/SnapshotManifest.java b/metadata/src/main/java/org/apache/kafka/image/loader/SnapshotManifest.java index b6c6dcce4d5ea..5653a4689ea12 100644 --- a/metadata/src/main/java/org/apache/kafka/image/loader/SnapshotManifest.java +++ b/metadata/src/main/java/org/apache/kafka/image/loader/SnapshotManifest.java @@ -25,7 +25,7 @@ /** * Contains information about a snapshot that was loaded. */ -public class SnapshotManifest { +public class SnapshotManifest implements LoaderManifest { /** * The source of this snapshot. */ @@ -44,6 +44,12 @@ public SnapshotManifest( this.elapsedNs = elapsedNs; } + @Override + public LoaderManifestType type() { + return LoaderManifestType.SNAPSHOT; + } + + @Override public MetadataProvenance provenance() { return provenance; } diff --git a/metadata/src/main/java/org/apache/kafka/image/publisher/MetadataPublisher.java b/metadata/src/main/java/org/apache/kafka/image/publisher/MetadataPublisher.java index 8dfba7a99abd8..46807fc2d364c 100644 --- a/metadata/src/main/java/org/apache/kafka/image/publisher/MetadataPublisher.java +++ b/metadata/src/main/java/org/apache/kafka/image/publisher/MetadataPublisher.java @@ -19,8 +19,8 @@ import org.apache.kafka.image.MetadataDelta; import org.apache.kafka.image.MetadataImage; -import org.apache.kafka.image.loader.LogDeltaManifest; -import org.apache.kafka.image.loader.SnapshotManifest; +import org.apache.kafka.image.loader.LoaderManifest; +import org.apache.kafka.raft.LeaderAndEpoch; /** @@ -40,33 +40,30 @@ public interface MetadataPublisher extends AutoCloseable { String name(); /** - * Publish a new cluster metadata snapshot that we loaded. + * Handle a change in the current controller. * - * @param delta The delta between the previous state and the new one. - * @param newImage The complete new state. - * @param manifest The contents of what was published. + * @param newLeaderAndEpoch The new quorum leader and epoch. The new leader will be + * OptionalInt.empty if there is currently no active controller. */ - void publishSnapshot( - MetadataDelta delta, - MetadataImage newImage, - SnapshotManifest manifest - ); + default void handleControllerChange(LeaderAndEpoch newLeaderAndEpoch) { } /** - * Publish a change to the cluster metadata. + * Publish a new cluster metadata snapshot that we loaded. * * @param delta The delta between the previous state and the new one. * @param newImage The complete new state. - * @param manifest The contents of what was published. + * @param manifest A manifest which describes the contents of what was published. + * If we loaded a snapshot, this will be a SnapshotManifest. + * If we loaded a log delta, this will be a LogDeltaManifest. */ - void publishLogDelta( + void publish( MetadataDelta delta, MetadataImage newImage, - LogDeltaManifest manifest + LoaderManifest manifest ); /** - * Close this metadata publisher. + * Close this metadata publisher and free any associated resources. */ - void close() throws Exception; + default void close() throws Exception { } } diff --git a/metadata/src/main/java/org/apache/kafka/image/publisher/SnapshotGenerator.java b/metadata/src/main/java/org/apache/kafka/image/publisher/SnapshotGenerator.java index 651acc5248334..dd795414aacb4 100644 --- a/metadata/src/main/java/org/apache/kafka/image/publisher/SnapshotGenerator.java +++ b/metadata/src/main/java/org/apache/kafka/image/publisher/SnapshotGenerator.java @@ -21,6 +21,7 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.image.MetadataDelta; import org.apache.kafka.image.MetadataImage; +import org.apache.kafka.image.loader.LoaderManifest; import org.apache.kafka.image.loader.LogDeltaManifest; import org.apache.kafka.image.loader.SnapshotManifest; import org.apache.kafka.queue.EventQueue; @@ -200,7 +201,24 @@ void resetSnapshotCounters() { } @Override - public void publishSnapshot( + public void publish( + MetadataDelta delta, + MetadataImage newImage, + LoaderManifest manifest + ) { + switch (manifest.type()) { + case LOG_DELTA: + publishLogDelta(delta, newImage, (LogDeltaManifest) manifest); + break; + case SNAPSHOT: + publishSnapshot(delta, newImage, (SnapshotManifest) manifest); + break; + default: + break; + } + } + + void publishSnapshot( MetadataDelta delta, MetadataImage newImage, SnapshotManifest manifest @@ -209,8 +227,7 @@ public void publishSnapshot( resetSnapshotCounters(); } - @Override - public void publishLogDelta( + void publishLogDelta( MetadataDelta delta, MetadataImage newImage, LogDeltaManifest manifest diff --git a/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java b/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java index 3394741f278c4..812404c3a7a73 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/migration/KRaftMigrationDriver.java @@ -23,8 +23,8 @@ import org.apache.kafka.image.MetadataDelta; import org.apache.kafka.image.MetadataImage; import org.apache.kafka.image.MetadataProvenance; -import org.apache.kafka.image.loader.LogDeltaManifest; -import org.apache.kafka.image.loader.SnapshotManifest; +import org.apache.kafka.image.loader.LoaderManifest; +import org.apache.kafka.image.loader.LoaderManifestType; import org.apache.kafka.image.publisher.MetadataPublisher; import org.apache.kafka.metadata.BrokerRegistration; import org.apache.kafka.queue.EventQueue; @@ -223,16 +223,21 @@ public String name() { } @Override - public void publishSnapshot(MetadataDelta delta, MetadataImage newImage, SnapshotManifest manifest) { - enqueueMetadataChangeEvent(delta, newImage, manifest.provenance(), true, NO_OP_HANDLER); + public void handleControllerChange(LeaderAndEpoch newLeaderAndEpoch) { + eventQueue.append(new KRaftLeaderEvent(newLeaderAndEpoch)); } @Override - public void publishLogDelta(MetadataDelta delta, MetadataImage newImage, LogDeltaManifest manifest) { - if (!leaderAndEpoch.equals(manifest.leaderAndEpoch())) { - eventQueue.append(new KRaftLeaderEvent(manifest.leaderAndEpoch())); - } - enqueueMetadataChangeEvent(delta, newImage, manifest.provenance(), false, NO_OP_HANDLER); + public void publish( + MetadataDelta delta, + MetadataImage newImage, + LoaderManifest manifest + ) { + enqueueMetadataChangeEvent(delta, + newImage, + manifest.provenance(), + manifest.type() == LoaderManifestType.SNAPSHOT, + NO_OP_HANDLER); } /** diff --git a/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java b/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java index b7e49243c43b3..71a806c3844cf 100644 --- a/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/loader/MetadataLoaderTest.java @@ -92,25 +92,23 @@ public String name() { } @Override - public void publishSnapshot( + public void publish( MetadataDelta delta, MetadataImage newImage, - SnapshotManifest manifest + LoaderManifest manifest ) { latestDelta = delta; latestImage = newImage; - latestSnapshotManifest = manifest; - } - - @Override - public void publishLogDelta( - MetadataDelta delta, - MetadataImage newImage, - LogDeltaManifest manifest - ) { - latestDelta = delta; - latestImage = newImage; - latestLogDeltaManifest = manifest; + switch (manifest.type()) { + case LOG_DELTA: + latestLogDeltaManifest = (LogDeltaManifest) manifest; + break; + case SNAPSHOT: + latestSnapshotManifest = (SnapshotManifest) manifest; + break; + default: + throw new RuntimeException("Invalid manifest type " + manifest.type()); + } } @Override diff --git a/metadata/src/test/java/org/apache/kafka/metadata/migration/KRaftMigrationDriverTest.java b/metadata/src/test/java/org/apache/kafka/metadata/migration/KRaftMigrationDriverTest.java index b25f1a10a765b..fd80ebce85bef 100644 --- a/metadata/src/test/java/org/apache/kafka/metadata/migration/KRaftMigrationDriverTest.java +++ b/metadata/src/test/java/org/apache/kafka/metadata/migration/KRaftMigrationDriverTest.java @@ -277,8 +277,9 @@ public void testOnlySendNeededRPCsToBrokers() throws Exception { image = delta.apply(provenance); // Publish a delta with this node (3000) as the leader - driver.publishLogDelta(delta, image, new LogDeltaManifest(provenance, - new LeaderAndEpoch(OptionalInt.of(3000), 1), 1, 100, 42)); + LeaderAndEpoch newLeader = new LeaderAndEpoch(OptionalInt.of(3000), 1); + driver.handleControllerChange(newLeader); + driver.publish(delta, image, new LogDeltaManifest(provenance, newLeader, 1, 100, 42)); TestUtils.waitForCondition(() -> driver.migrationState().get(1, TimeUnit.MINUTES).equals(MigrationDriverState.DUAL_WRITE), "Waiting for KRaftMigrationDriver to enter DUAL_WRITE state");