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..d8448935e0a4e 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.