From 44eea04f4a7cf2a2a1a7016bcb23da0511d6cc38 Mon Sep 17 00:00:00 2001 From: Leonard Ge Date: Tue, 21 Apr 2020 09:47:55 +0100 Subject: [PATCH 1/6] Avoid starting election for topics where preferred leader is not in sync. --- core/src/main/scala/kafka/controller/KafkaController.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index 9ad7b6ff7b634..c08f7cb378820 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1066,6 +1066,7 @@ class KafkaController(val config: KafkaConfig, // do this check only if the broker is live and there are no partitions being reassigned currently // and preferred replica election is not in progress val candidatePartitions = topicsNotInPreferredReplica.keys.filter(tp => controllerContext.isReplicaOnline(leaderBroker, tp) && + controllerContext.partitionLeadershipInfo(tp).leaderAndIsr.isr.contains(leaderBroker) && controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic)) From 29632f76ed34c75f80c175101c17155e98f24f0a Mon Sep 17 00:00:00 2001 From: Leonard Ge Date: Wed, 22 Apr 2020 11:13:24 +0100 Subject: [PATCH 2/6] Refactored the code and fixed bug in the integration test. --- core/src/main/scala/kafka/controller/KafkaController.scala | 5 +++-- .../unit/kafka/controller/ControllerIntegrationTest.scala | 4 ++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index c08f7cb378820..c82e5e9a6034a 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1066,10 +1066,11 @@ class KafkaController(val config: KafkaConfig, // do this check only if the broker is live and there are no partitions being reassigned currently // and preferred replica election is not in progress val candidatePartitions = topicsNotInPreferredReplica.keys.filter(tp => controllerContext.isReplicaOnline(leaderBroker, tp) && - controllerContext.partitionLeadershipInfo(tp).leaderAndIsr.isr.contains(leaderBroker) && controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && - controllerContext.allTopics.contains(tp.topic)) + controllerContext.allTopics.contains(tp.topic) && + controllerContext.partitionLeadershipInfo.get(tp).forall(l => l.leaderAndIsr.isr.contains(leaderBroker)) + ) onReplicaElection(candidatePartitions.toSet, ElectionType.PREFERRED, AutoTriggered) } } diff --git a/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala b/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala index c4b5f478d1953..c7a1cd5b840fb 100644 --- a/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ControllerIntegrationTest.scala @@ -433,8 +433,8 @@ class ControllerIntegrationTest extends ZooKeeperTestHarness { TestUtils.createTopic(zkClient, tp.topic, partitionReplicaAssignment = assignment, servers = servers) waitForPartitionState(tp, firstControllerEpoch, otherBrokerId, LeaderAndIsr.initialLeaderEpoch, "failed to get expected partition state upon topic creation") - servers(1).shutdown() - servers(1).awaitShutdown() + servers(otherBrokerId).shutdown() + servers(otherBrokerId).awaitShutdown() TestUtils.waitUntilTrue(() => { val leaderIsrAndControllerEpochMap = zkClient.getTopicPartitionStates(Seq(tp)) leaderIsrAndControllerEpochMap.contains(tp) && From 172960f972e82b10158d86fb4c02a7f9317ac43c Mon Sep 17 00:00:00 2001 From: Leonard Ge Date: Thu, 23 Apr 2020 10:03:04 +0100 Subject: [PATCH 3/6] Refactored according to comments. --- core/src/main/scala/kafka/controller/KafkaController.scala | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index c82e5e9a6034a..fb0c6aeb26a82 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1069,8 +1069,11 @@ class KafkaController(val config: KafkaConfig, controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic) && - controllerContext.partitionLeadershipInfo.get(tp).forall(l => l.leaderAndIsr.isr.contains(leaderBroker)) - ) + PartitionLeaderElectionAlgorithms.preferredReplicaPartitionLeaderElection( + controllerContext.partitionReplicaAssignment(tp), + controllerContext.partitionLeadershipInfo(tp).leaderAndIsr.isr, + controllerContext.liveBrokerIds.toSet).nonEmpty + ) onReplicaElection(candidatePartitions.toSet, ElectionType.PREFERRED, AutoTriggered) } } From 24c6d42b9ea69fe9bb2a40649bae15e0f74542db Mon Sep 17 00:00:00 2001 From: Leonard Ge Date: Fri, 24 Apr 2020 10:08:02 +0100 Subject: [PATCH 4/6] Kept the code to be consistent. --- .../kafka/controller/KafkaController.scala | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index fb0c6aeb26a82..fcae79247472d 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1065,20 +1065,27 @@ class KafkaController(val config: KafkaConfig, if (imbalanceRatio > (config.leaderImbalancePerBrokerPercentage.toDouble / 100)) { // do this check only if the broker is live and there are no partitions being reassigned currently // and preferred replica election is not in progress - val candidatePartitions = topicsNotInPreferredReplica.keys.filter(tp => controllerContext.isReplicaOnline(leaderBroker, tp) && + val candidatePartitions = topicsNotInPreferredReplica.keys.filter(tp => + controllerContext.isReplicaOnline(leaderBroker, tp) && controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic) && - PartitionLeaderElectionAlgorithms.preferredReplicaPartitionLeaderElection( - controllerContext.partitionReplicaAssignment(tp), - controllerContext.partitionLeadershipInfo(tp).leaderAndIsr.isr, - controllerContext.liveBrokerIds.toSet).nonEmpty + isPreferredLeaderInSync(tp) ) onReplicaElection(candidatePartitions.toSet, ElectionType.PREFERRED, AutoTriggered) } } } + private def isPreferredLeaderInSync(tp: TopicPartition): Boolean = { + val assignment = controllerContext.partitionReplicaAssignment(tp) + val liveReplicas = assignment.filter(replica => controllerContext.isReplicaOnline(replica, tp)) + val isr = controllerContext.partitionLeadershipInfo(tp).leaderAndIsr.isr + PartitionLeaderElectionAlgorithms + .preferredReplicaPartitionLeaderElection(assignment, isr, liveReplicas.toSet) + .nonEmpty + } + private def processAutoPreferredReplicaLeaderElection(): Unit = { if (!isActive) return try { From dc67b99c2d2d45a59ede67da186d1890c7b7c1ce Mon Sep 17 00:00:00 2001 From: Leonard Ge Date: Fri, 24 Apr 2020 10:09:25 +0100 Subject: [PATCH 5/6] Reverted style change. --- core/src/main/scala/kafka/controller/KafkaController.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index fcae79247472d..ad5c89892f8ff 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1065,8 +1065,7 @@ class KafkaController(val config: KafkaConfig, if (imbalanceRatio > (config.leaderImbalancePerBrokerPercentage.toDouble / 100)) { // do this check only if the broker is live and there are no partitions being reassigned currently // and preferred replica election is not in progress - val candidatePartitions = topicsNotInPreferredReplica.keys.filter(tp => - controllerContext.isReplicaOnline(leaderBroker, tp) && + val candidatePartitions = topicsNotInPreferredReplica.keys.filter(tp => controllerContext.isReplicaOnline(leaderBroker, tp) && controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic) && From 2fa3ceed0322b90306393e0dffc1fdbdba779613 Mon Sep 17 00:00:00 2001 From: Leonard Ge <62600326+leonardge@users.noreply.github.com> Date: Sat, 25 Apr 2020 16:25:39 +0100 Subject: [PATCH 6/6] Renamed helper function after review. --- core/src/main/scala/kafka/controller/KafkaController.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index ad5c89892f8ff..82be66ae2cbe2 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1069,14 +1069,14 @@ class KafkaController(val config: KafkaConfig, controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic) && - isPreferredLeaderInSync(tp) + canPreferredReplicaBeLeader(tp) ) onReplicaElection(candidatePartitions.toSet, ElectionType.PREFERRED, AutoTriggered) } } } - private def isPreferredLeaderInSync(tp: TopicPartition): Boolean = { + private def canPreferredReplicaBeLeader(tp: TopicPartition): Boolean = { val assignment = controllerContext.partitionReplicaAssignment(tp) val liveReplicas = assignment.filter(replica => controllerContext.isReplicaOnline(replica, tp)) val isr = controllerContext.partitionLeadershipInfo(tp).leaderAndIsr.isr