From 07fa570b9a8254c55990513ae1ca1dfe2aa739b7 Mon Sep 17 00:00:00 2001 From: Lucas Wang Date: Tue, 7 Sep 2021 14:32:47 -0700 Subject: [PATCH 1/2] [LI-HOTFIX] Improving performance of the isReplicaOnline method TICKET = LIKAFKA-38536 LI_DESCRIPTION = The ControllerContext.isReplicaOnline method is frequently called inside a loop, and internally this method relies on several derived fields in the ControllerContext. Repeatedly calculating the derived fields could be expensive, and yet these fields do not change between iterations of the loop. This PR creates a new class ControllerContextSnapshot that tries to cache the derived fields in order to save the repeated computation and speed up the controller. EXIT_CRITERIA = When this change is proposed in upstream kafka and pulled internally. --- .../controller/ControllerChannelManager.scala | 4 +-- .../kafka/controller/ControllerContext.scala | 33 +++++++++++++------ .../scala/kafka/controller/Election.scala | 15 +++++---- .../kafka/controller/KafkaController.scala | 13 +++++--- .../controller/PartitionStateMachine.scala | 9 +++-- .../controller/ReplicaStateMachine.scala | 3 +- 6 files changed, 50 insertions(+), 27 deletions(-) diff --git a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala index 9b4ffc4b27187..40658cd1f9abb 100755 --- a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala @@ -508,12 +508,12 @@ abstract class AbstractControllerBrokerRequestBatch(config: KafkaConfig, /** Send UpdateMetadataRequest to the given brokers for the given partitions and partitions that are being deleted */ def addUpdateMetadataRequestForBrokers(brokerIds: Seq[Int], partitions: collection.Set[TopicPartition]): Unit = { - + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) def updateMetadataRequestPartitionInfo(partition: TopicPartition, beingDeleted: Boolean): Unit = { controllerContext.partitionLeadershipInfo.get(partition) match { case Some(LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) => val replicas = controllerContext.partitionReplicaAssignment(partition) - val offlineReplicas = replicas.filter(!controllerContext.isReplicaOnline(_, partition)) + val offlineReplicas = replicas.filter(!controllerContextSnapshot.isReplicaOnline(_, partition)) val updatedLeaderAndIsr = if (beingDeleted) LeaderAndIsr.duringDelete(leaderAndIsr.isr) else leaderAndIsr diff --git a/core/src/main/scala/kafka/controller/ControllerContext.scala b/core/src/main/scala/kafka/controller/ControllerContext.scala index 5f70f74c91bb7..59b97d44a7314 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -225,14 +225,6 @@ class ControllerContext { }.toSet } - def isReplicaOnline(brokerId: Int, topicPartition: TopicPartition, includeShuttingDownBrokers: Boolean = false): Boolean = { - val brokerOnline = { - if (includeShuttingDownBrokers) liveOrShuttingDownBrokerIds.contains(brokerId) - else liveBrokerIds.contains(brokerId) - } - brokerOnline && !replicasOnOfflineDirs.getOrElse(brokerId, Set.empty).contains(topicPartition) - } - def replicasOnBrokers(brokerIds: Set[Int]): Set[PartitionAndReplica] = { brokerIds.flatMap { brokerId => partitionAssignments.flatMap { @@ -257,8 +249,9 @@ class ControllerContext { } def allLiveReplicas(): Set[PartitionAndReplica] = { + val snapshot = ControllerContextSnapshot(this) replicasOnBrokers(liveBrokerIds).filter { partitionAndReplica => - isReplicaOnline(partitionAndReplica.replica, partitionAndReplica.topicPartition) + snapshot.isReplicaOnline(partitionAndReplica.replica, partitionAndReplica.topicPartition) } } @@ -270,12 +263,13 @@ class ControllerContext { def onlineAndOfflineReplicas: (Set[PartitionAndReplica], Set[PartitionAndReplica]) = { val onlineReplicas = mutable.Set.empty[PartitionAndReplica] val offlineReplicas = mutable.Set.empty[PartitionAndReplica] + val snapshot = ControllerContextSnapshot(this) for ((topic, partitionAssignments) <- partitionAssignments; (partitionId, assignment) <- partitionAssignments) { val partition = new TopicPartition(topic, partitionId) for (replica <- assignment.replicas) { val partitionAndReplica = PartitionAndReplica(partition, replica) - if (isReplicaOnline(replica, partition)) + if (snapshot.isReplicaOnline(replica, partition)) onlineReplicas.add(partitionAndReplica) else offlineReplicas.add(partitionAndReplica) @@ -431,3 +425,22 @@ class ControllerContext { targetState.validPreviousStates.contains(partitionStates(partition)) } + +/** + * The ControllerContextSnapshot is an immutable snapshot of the ControllorContext. + * The motivation for this class is that we don't need to calculate certain fields + * repeatedly, including liveBrokerIds and liveOrShuttingDownBrokerIds. + */ +case class ControllerContextSnapshot(controllerContext: ControllerContext) { + val liveBrokerIds = controllerContext.liveBrokerIds + val liveOrShuttingDownBrokerIds = controllerContext.liveOrShuttingDownBrokerIds + val replicasOnOfflineDirs = controllerContext.replicasOnOfflineDirs + + def isReplicaOnline(brokerId: Int, topicPartition: TopicPartition, includeShuttingDownBrokers: Boolean = false): Boolean = { + val brokerOnline = { + if (includeShuttingDownBrokers) liveOrShuttingDownBrokerIds.contains(brokerId) + else liveBrokerIds.contains(brokerId) + } + brokerOnline && !replicasOnOfflineDirs.getOrElse(brokerId, Set.empty).contains(topicPartition) + } +} diff --git a/core/src/main/scala/kafka/controller/Election.scala b/core/src/main/scala/kafka/controller/Election.scala index ae6c32a82f030..3b6d8daf89054 100644 --- a/core/src/main/scala/kafka/controller/Election.scala +++ b/core/src/main/scala/kafka/controller/Election.scala @@ -30,9 +30,9 @@ object Election extends Logging { leaderAndIsrOpt: Option[LeaderAndIsr], uncleanLeaderElectionEnabled: Boolean, controllerContext: ControllerContext): ElectionResult = { - + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val assignment = controllerContext.partitionReplicaAssignment(partition) - val liveReplicas = assignment.filter(replica => controllerContext.isReplicaOnline(replica, partition)) + val liveReplicas = assignment.filter(replica => controllerContextSnapshot.isReplicaOnline(replica, partition)) leaderAndIsrOpt match { case Some(leaderAndIsr) => val isr = leaderAndIsr.isr @@ -44,7 +44,7 @@ object Election extends Logging { warn(s"Unclean leader election. Partition $partition has been assigned leader $leader from deposed " + s"leader ${leaderAndIsr.leader}.") } - val newIsr = if (isr.contains(leader)) isr.filter(replica => controllerContext.isReplicaOnline(replica, partition)) + val newIsr = if (isr.contains(leader)) isr.filter(replica => controllerContextSnapshot.isReplicaOnline(replica, partition)) else List(leader) leaderAndIsr.newLeaderAndIsr(leader, newIsr) } @@ -78,8 +78,9 @@ object Election extends Logging { private def leaderForReassign(partition: TopicPartition, leaderAndIsr: LeaderAndIsr, controllerContext: ControllerContext): ElectionResult = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val targetReplicas = controllerContext.partitionFullReplicaAssignment(partition).targetReplicas - val liveReplicas = targetReplicas.filter(replica => controllerContext.isReplicaOnline(replica, partition)) + val liveReplicas = targetReplicas.filter(replica => controllerContextSnapshot.isReplicaOnline(replica, partition)) val isr = leaderAndIsr.isr val leaderOpt = PartitionLeaderElectionAlgorithms.reassignPartitionLeaderElection(targetReplicas, isr, liveReplicas.toSet) val newLeaderAndIsrOpt = leaderOpt.map(leader => leaderAndIsr.newLeader(leader)) @@ -105,8 +106,9 @@ object Election extends Logging { private def leaderForPreferredReplica(partition: TopicPartition, leaderAndIsr: LeaderAndIsr, controllerContext: ControllerContext): ElectionResult = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val assignment = controllerContext.partitionReplicaAssignment(partition) - val liveReplicas = assignment.filter(replica => controllerContext.isReplicaOnline(replica, partition)) + val liveReplicas = assignment.filter(replica => controllerContextSnapshot.isReplicaOnline(replica, partition)) val isr = leaderAndIsr.isr val leaderOpt = PartitionLeaderElectionAlgorithms.preferredReplicaPartitionLeaderElection(assignment, isr, liveReplicas.toSet) val newLeaderAndIsrOpt = leaderOpt.map(leader => leaderAndIsr.newLeader(leader)) @@ -133,9 +135,10 @@ object Election extends Logging { leaderAndIsr: LeaderAndIsr, shuttingDownBrokerIds: Set[Int], controllerContext: ControllerContext): ElectionResult = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val assignment = controllerContext.partitionReplicaAssignment(partition) val liveOrShuttingDownReplicas = assignment.filter(replica => - controllerContext.isReplicaOnline(replica, partition, includeShuttingDownBrokers = true)) + controllerContextSnapshot.isReplicaOnline(replica, partition, includeShuttingDownBrokers = true)) val isr = leaderAndIsr.isr val leaderOpt = PartitionLeaderElectionAlgorithms.controlledShutdownPartitionLeaderElection(assignment, isr, liveOrShuttingDownReplicas.toSet, shuttingDownBrokerIds) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index 356a5ec96d8bb..0dd27c97b5f1a 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -571,11 +571,12 @@ class KafkaController(val config: KafkaConfig, * partitions coming online. */ private def onReplicasBecomeOffline(newOfflineReplicas: Set[PartitionAndReplica]): Unit = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val (newOfflineReplicasForDeletion, newOfflineReplicasNotForDeletion) = newOfflineReplicas.partition(p => topicDeletionManager.isTopicQueuedUpForDeletion(p.topic)) val partitionsWithoutLeader = controllerContext.partitionLeadershipInfo.filter(partitionAndLeader => - !controllerContext.isReplicaOnline(partitionAndLeader._2.leaderAndIsr.leader, partitionAndLeader._1) && + !controllerContextSnapshot.isReplicaOnline(partitionAndLeader._2.leaderAndIsr.leader, partitionAndLeader._1) && !topicDeletionManager.isTopicQueuedUpForDeletion(partitionAndLeader._1.topic)).keySet // trigger OfflinePartition state for all partitions whose current leader is one amongst the newOfflineReplicas @@ -935,10 +936,11 @@ class KafkaController(val config: KafkaConfig, } private def fetchTopicDeletionsInProgress(): (Set[String], Set[String]) = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val topicsToBeDeleted = zkClient.getTopicDeletions.toSet val topicsWithOfflineReplicas = controllerContext.allTopics.filter { topic => { val replicasForTopic = controllerContext.replicasForTopic(topic) - replicasForTopic.exists(r => !controllerContext.isReplicaOnline(r.replica, r.topicPartition)) + replicasForTopic.exists(r => !controllerContextSnapshot.isReplicaOnline(r.replica, r.topicPartition)) }} val topicsForWhichPartitionReassignmentIsInProgress = controllerContext.partitionsBeingReassigned.map(_.topic) val topicsIneligibleForDeletion = topicsWithOfflineReplicas | topicsForWhichPartitionReassignmentIsInProgress @@ -1001,13 +1003,13 @@ class KafkaController(val config: KafkaConfig, newAssignment: ReplicaAssignment): Unit = { val reassignedReplicas = newAssignment.replicas val currentLeader = controllerContext.partitionLeadershipInfo(topicPartition).leaderAndIsr.leader - + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) if (!reassignedReplicas.contains(currentLeader)) { info(s"Leader $currentLeader for partition $topicPartition being reassigned, " + s"is not in the new list of replicas ${reassignedReplicas.mkString(",")}. Re-electing leader") // move the leader to one of the alive and caught up new replicas partitionStateMachine.handleStateChanges(Seq(topicPartition), OnlinePartition, Some(ReassignPartitionLeaderElectionStrategy)) - } else if (controllerContext.isReplicaOnline(currentLeader, topicPartition)) { + } else if (controllerContextSnapshot.isReplicaOnline(currentLeader, topicPartition)) { info(s"Leader $currentLeader for partition $topicPartition being reassigned, " + s"is already in the new list of replicas ${reassignedReplicas.mkString(",")} and is alive") // shrink replication factor and update the leader epoch in zookeeper to use on the next LeaderAndIsrRequest @@ -1225,6 +1227,7 @@ class KafkaController(val config: KafkaConfig, private def checkAndTriggerAutoLeaderRebalance(): Unit = { trace("Checking need to trigger auto leader balancing") + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val preferredReplicasForTopicsByBrokers: Map[Int, Map[TopicPartition, Seq[Int]]] = controllerContext.allPartitions.filterNot { tp => topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) @@ -1250,7 +1253,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 => controllerContextSnapshot.isReplicaOnline(leaderBroker, tp) && controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic)) diff --git a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala index b1ecedfcd0684..a027572ec44a8 100755 --- a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala +++ b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala @@ -80,12 +80,13 @@ abstract class PartitionStateMachine(controllerContext: ControllerContext) exten * zookeeper */ private def initializePartitionState(): Unit = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) for (topicPartition <- controllerContext.allPartitions) { // check if leader and isr path exists for partition. If not, then it is in NEW state controllerContext.partitionLeadershipInfo.get(topicPartition) match { case Some(currentLeaderIsrAndEpoch) => // else, check if the leader for partition is alive. If yes, it is in Online state, else it is in Offline state - if (controllerContext.isReplicaOnline(currentLeaderIsrAndEpoch.leaderAndIsr.leader, topicPartition)) + if (controllerContextSnapshot.isReplicaOnline(currentLeaderIsrAndEpoch.leaderAndIsr.leader, topicPartition)) // leader is alive controllerContext.putPartitionState(topicPartition, OnlinePartition) else @@ -271,10 +272,11 @@ class ZkPartitionStateMachine(config: KafkaConfig, * @return The partitions that have been successfully initialized. */ private def initializeLeaderAndIsrForPartitions(partitions: Seq[TopicPartition]): Seq[TopicPartition] = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val successfulInitializations = mutable.Buffer.empty[TopicPartition] val replicasPerPartition = partitions.map(partition => partition -> controllerContext.partitionReplicaAssignment(partition)) val liveReplicasPerPartition = replicasPerPartition.map { case (partition, replicas) => - val liveReplicasForPartition = replicas.filter(replica => controllerContext.isReplicaOnline(replica, partition)) + val liveReplicasForPartition = replicas.filter(replica => controllerContextSnapshot.isReplicaOnline(replica, partition)) partition -> liveReplicasForPartition } val (partitionsWithoutLiveReplicas, partitionsWithLiveReplicas) = liveReplicasPerPartition.partition { case (_, liveReplicas) => liveReplicas.isEmpty } @@ -460,9 +462,10 @@ class ZkPartitionStateMachine(config: KafkaConfig, leaderAndIsrs: Seq[(TopicPartition, LeaderAndIsr)], allowUnclean: Boolean ): Seq[(TopicPartition, Option[LeaderAndIsr], Boolean)] = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) val (partitionsWithNoLiveInSyncReplicas, partitionsWithLiveInSyncReplicas) = leaderAndIsrs.partition { case (partition, leaderAndIsr) => - val liveInSyncReplicas = leaderAndIsr.isr.filter(controllerContext.isReplicaOnline(_, partition)) + val liveInSyncReplicas = leaderAndIsr.isr.filter(controllerContextSnapshot.isReplicaOnline(_, partition)) liveInSyncReplicas.isEmpty } diff --git a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala index 88c49f832464c..83fead8717750 100644 --- a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala +++ b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala @@ -55,11 +55,12 @@ abstract class ReplicaStateMachine(controllerContext: ControllerContext) extends * in zookeeper */ private def initializeReplicaState(): Unit = { + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) controllerContext.allPartitions.foreach { partition => val replicas = controllerContext.partitionReplicaAssignment(partition) replicas.foreach { replicaId => val partitionAndReplica = PartitionAndReplica(partition, replicaId) - if (controllerContext.isReplicaOnline(replicaId, partition)) { + if (controllerContextSnapshot.isReplicaOnline(replicaId, partition)) { controllerContext.putReplicaState(partitionAndReplica, OnlineReplica) } else { // mark replicas on dead brokers as failed for topic deletion, if they belong to a topic to be deleted. From b5cfbd55c167bfa528557bbe7095cf7b3b9c4c85 Mon Sep 17 00:00:00 2001 From: Lucas Wang Date: Thu, 9 Sep 2021 14:12:33 -0700 Subject: [PATCH 2/2] Fixing style --- .../main/scala/kafka/controller/ControllerContext.scala | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/controller/ControllerContext.scala b/core/src/main/scala/kafka/controller/ControllerContext.scala index 59b97d44a7314..d8448935e0a4e 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -427,10 +427,10 @@ class ControllerContext { } /** - * The ControllerContextSnapshot is an immutable snapshot of the ControllorContext. - * The motivation for this class is that we don't need to calculate certain fields - * repeatedly, including liveBrokerIds and liveOrShuttingDownBrokerIds. - */ + * The ControllerContextSnapshot is an immutable snapshot of the ControllorContext. + * The motivation for this class is that we don't need to calculate certain fields + * repeatedly, including liveBrokerIds and liveOrShuttingDownBrokerIds. + */ case class ControllerContextSnapshot(controllerContext: ControllerContext) { val liveBrokerIds = controllerContext.liveBrokerIds val liveOrShuttingDownBrokerIds = controllerContext.liveOrShuttingDownBrokerIds