From f08c91ef1cd83a31b28a0bd71428c9446974143b Mon Sep 17 00:00:00 2001 From: David Arthur Date: Tue, 18 Apr 2023 14:41:33 -0400 Subject: [PATCH 1/3] Only send ZK RPCs to unfenced ZK brokers --- .../kafka/migration/MigrationPropagator.scala | 26 +++++++++++++------ 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/migration/MigrationPropagator.scala b/core/src/main/scala/kafka/migration/MigrationPropagator.scala index 80f462bbd21c1..5768ce18451ef 100644 --- a/core/src/main/scala/kafka/migration/MigrationPropagator.scala +++ b/core/src/main/scala/kafka/migration/MigrationPropagator.scala @@ -70,14 +70,24 @@ class MigrationPropagator( override def publishMetadata(image: MetadataImage): Unit = { val oldImage = _image - val addedBrokers = new util.HashSet[Integer](image.cluster().brokers().keySet()) - addedBrokers.removeAll(oldImage.cluster().brokers().keySet()) - val removedBrokers = new util.HashSet[Integer](oldImage.cluster().brokers().keySet()) - removedBrokers.removeAll(image.cluster().brokers().keySet()) - - removedBrokers.asScala.foreach(id => channelManager.removeBroker(id)) - addedBrokers.asScala.foreach(id => - channelManager.addBroker(Broker.fromBrokerRegistration(image.cluster().broker(id)))) + val prevBrokers = oldImage.cluster().brokers().values().asScala + .filter(_.isMigratingZkBroker) + .filterNot(_.fenced) + .map(Broker.fromBrokerRegistration) + .toSet + + val aliveBrokers = image.cluster().brokers().values().asScala + .filter(_.isMigratingZkBroker) + .filterNot(_.fenced) + .map(Broker.fromBrokerRegistration) + .toSet + + val addedBrokers = aliveBrokers -- prevBrokers + val removedBrokers = prevBrokers -- aliveBrokers + + stateChangeLogger.logger.debug(s"Adding brokers $addedBrokers, removing brokers $removedBrokers.") + removedBrokers.foreach(broker => channelManager.removeBroker(broker.id)) + addedBrokers.foreach(broker => channelManager.addBroker(broker)) _image = image } From 4b9aa7c8f5c91fc57e8bbb308a6b72cbafe10e09 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Tue, 18 Apr 2023 15:43:33 -0400 Subject: [PATCH 2/3] Add unit test for calculateBrokerChanges --- .../kafka/migration/MigrationPropagator.scala | 41 +++++++---- .../migration/MigrationPropagatorTest.scala | 73 +++++++++++++++++++ 2 files changed, 98 insertions(+), 16 deletions(-) create mode 100644 core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala diff --git a/core/src/main/scala/kafka/migration/MigrationPropagator.scala b/core/src/main/scala/kafka/migration/MigrationPropagator.scala index 5768ce18451ef..e7bbd11edd9ef 100644 --- a/core/src/main/scala/kafka/migration/MigrationPropagator.scala +++ b/core/src/main/scala/kafka/migration/MigrationPropagator.scala @@ -23,7 +23,7 @@ import kafka.server.KafkaConfig import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.metrics.Metrics import org.apache.kafka.common.utils.Time -import org.apache.kafka.image.{MetadataDelta, MetadataImage, TopicsImage} +import org.apache.kafka.image.{ClusterImage, MetadataDelta, MetadataImage, TopicsImage} import org.apache.kafka.metadata.PartitionRegistration import org.apache.kafka.metadata.migration.LegacyPropagator import org.apache.kafka.server.common.MetadataVersion @@ -32,6 +32,26 @@ import java.util import scala.jdk.CollectionConverters._ import scala.compat.java8.OptionConverters._ +object MigrationPropagator { + def calculateBrokerChanges(prevClusterImage: ClusterImage, clusterImage: ClusterImage): (Set[Broker], Set[Broker]) = { + val prevBrokers = prevClusterImage.brokers().values().asScala + .filter(_.isMigratingZkBroker) + .filterNot(_.fenced) + .map(Broker.fromBrokerRegistration) + .toSet + + val aliveBrokers = clusterImage.brokers().values().asScala + .filter(_.isMigratingZkBroker) + .filterNot(_.fenced) + .map(Broker.fromBrokerRegistration) + .toSet + + val addedBrokers = aliveBrokers -- prevBrokers + val removedBrokers = prevBrokers -- aliveBrokers + (addedBrokers, removedBrokers) + } +} + class MigrationPropagator( nodeId: Int, config: KafkaConfig @@ -70,22 +90,11 @@ class MigrationPropagator( override def publishMetadata(image: MetadataImage): Unit = { val oldImage = _image - val prevBrokers = oldImage.cluster().brokers().values().asScala - .filter(_.isMigratingZkBroker) - .filterNot(_.fenced) - .map(Broker.fromBrokerRegistration) - .toSet - - val aliveBrokers = image.cluster().brokers().values().asScala - .filter(_.isMigratingZkBroker) - .filterNot(_.fenced) - .map(Broker.fromBrokerRegistration) - .toSet - val addedBrokers = aliveBrokers -- prevBrokers - val removedBrokers = prevBrokers -- aliveBrokers - - stateChangeLogger.logger.debug(s"Adding brokers $addedBrokers, removing brokers $removedBrokers.") + val (addedBrokers, removedBrokers) = MigrationPropagator.calculateBrokerChanges(oldImage.cluster(), image.cluster()) + if (addedBrokers.nonEmpty || removedBrokers.nonEmpty) { + stateChangeLogger.logger.info(s"Adding brokers $addedBrokers, removing brokers $removedBrokers.") + } removedBrokers.foreach(broker => channelManager.removeBroker(broker.id)) addedBrokers.foreach(broker => channelManager.addBroker(broker)) _image = image diff --git a/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala b/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala new file mode 100644 index 0000000000000..29e5870450b2c --- /dev/null +++ b/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala @@ -0,0 +1,73 @@ +package kafka.migration + +import kafka.cluster.Broker +import org.apache.kafka.common.metadata.RegisterBrokerRecord +import org.apache.kafka.image.ClusterImage +import org.apache.kafka.metadata.BrokerRegistration +import org.junit.jupiter.api.Assertions.{assertFalse, assertTrue} +import org.junit.jupiter.api.Test + +import scala.jdk.CollectionConverters._ + +class MigrationPropagatorTest { + def brokerBuilder(brokerId: Int, isZkBroker: Boolean, isFenced: Boolean): BrokerRegistration = { + BrokerRegistration.fromRecord( + new RegisterBrokerRecord() + .setBrokerId(brokerId) + .setIsMigratingZkBroker(isZkBroker) + .setBrokerEpoch(10) + .setFenced(isFenced) + ) + } + + def brokersToClusterImage(brokers: Seq[BrokerRegistration]): ClusterImage = { + val brokerMap = brokers.map(broker => Integer.valueOf(broker.id()) -> broker).toMap.asJava + new ClusterImage(brokerMap) + } + + @Test + def testCalculateBrokerChanges(): Unit = { + // Start with one fenced, one un-fenced ZK broker + var broker0 = brokerBuilder(0, true, true) + var broker1 = brokerBuilder(1, true, false) + MigrationPropagator.calculateBrokerChanges(ClusterImage.EMPTY, brokersToClusterImage(Seq(broker0, broker1))) match { + case (addedBrokers, removedBrokers) => + assertFalse(addedBrokers.contains(Broker.fromBrokerRegistration(broker0))) + assertTrue(addedBrokers.contains(Broker.fromBrokerRegistration(broker1))) + assertTrue(removedBrokers.isEmpty) + } + + // Un-fence broker 0 + var prevImage = brokersToClusterImage(Seq(broker0, broker1)) + broker0 = brokerBuilder(0, true, false) + broker1 = brokerBuilder(1, true, false) + MigrationPropagator.calculateBrokerChanges(prevImage, brokersToClusterImage(Seq(broker0, broker1))) match { + case (addedBrokers, removedBrokers) => + assertTrue(addedBrokers.contains(Broker.fromBrokerRegistration(broker0))) + assertFalse(addedBrokers.contains(Broker.fromBrokerRegistration(broker1))) + assertTrue(removedBrokers.isEmpty) + } + + // Migrate both to KRaft + prevImage = brokersToClusterImage(Seq(broker0, broker1)) + broker0 = brokerBuilder(0, false, false) + broker1 = brokerBuilder(1, false, false) + MigrationPropagator.calculateBrokerChanges(prevImage, brokersToClusterImage(Seq(broker0, broker1))) match { + case (addedBrokers, removedBrokers) => + assertTrue(addedBrokers.isEmpty) + assertTrue(removedBrokers.contains(Broker.fromBrokerRegistration(broker0))) + assertTrue(removedBrokers.contains(Broker.fromBrokerRegistration(broker0))) + } + + // Downgrade one back to ZK + prevImage = brokersToClusterImage(Seq(broker0, broker1)) + broker0 = brokerBuilder(0, true, false) + broker1 = brokerBuilder(1, false, false) + MigrationPropagator.calculateBrokerChanges(prevImage, brokersToClusterImage(Seq(broker0, broker1))) match { + case (addedBrokers, removedBrokers) => + assertTrue(addedBrokers.contains(Broker.fromBrokerRegistration(broker0))) + assertFalse(addedBrokers.contains(Broker.fromBrokerRegistration(broker1))) + assertTrue(removedBrokers.isEmpty) + } + } +} From 645140f9000662609a41f945e0673b1d353f60bd Mon Sep 17 00:00:00 2001 From: David Arthur Date: Tue, 18 Apr 2023 15:53:27 -0400 Subject: [PATCH 3/3] add license --- .../migration/MigrationPropagatorTest.scala | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala b/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala index 29e5870450b2c..b7cdb57cc887e 100644 --- a/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala +++ b/core/src/test/scala/unit/kafka/migration/MigrationPropagatorTest.scala @@ -1,3 +1,20 @@ +/* + * 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 kafka.migration import kafka.cluster.Broker