diff --git a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala index f5f1555fe3ac2..74ed53a626ecd 100755 --- a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala @@ -441,12 +441,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(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 28a6eef539aab..df3eb3de5c24a 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -238,14 +238,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 { @@ -279,12 +271,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) @@ -469,8 +462,9 @@ class ControllerContext { partitionLeadershipInfo.keySet.filter(tp => !isTopicQueuedUpForDeletion(tp.topic)) def partitionsWithOfflineLeader: Set[TopicPartition] = { + val snapshot = ControllerContextSnapshot(this) partitionLeadershipInfo.filter { case (topicPartition, leaderIsrAndControllerEpoch) => - !isReplicaOnline(leaderIsrAndControllerEpoch.leaderAndIsr.leader, topicPartition) && + !snapshot.isReplicaOnline(leaderIsrAndControllerEpoch.leaderAndIsr.leader, topicPartition) && !isTopicQueuedUpForDeletion(topicPartition.topic) }.keySet } @@ -539,3 +533,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 ab7c200327734..8a747279d455f 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1043,10 +1043,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 @@ -1109,13 +1110,14 @@ class KafkaController(val config: KafkaConfig, newAssignment: ReplicaAssignment): Unit = { val reassignedReplicas = newAssignment.replicas val currentLeader = controllerContext.partitionLeadershipInfo(topicPartition).get.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 @@ -1347,6 +1349,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) @@ -1374,16 +1377,16 @@ class KafkaController(val config: KafkaConfig, controllerContext.partitionsBeingReassigned.isEmpty && !topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) && controllerContext.allTopics.contains(tp.topic) && - canPreferredReplicaBeLeader(tp) + canPreferredReplicaBeLeader(tp, controllerContextSnapshot) ) onReplicaElection(candidatePartitions.toSet, ElectionType.PREFERRED, AutoTriggered) } } } - private def canPreferredReplicaBeLeader(tp: TopicPartition): Boolean = { + private def canPreferredReplicaBeLeader(tp: TopicPartition, controllerContextSnapshot: ControllerContextSnapshot): Boolean = { val assignment = controllerContext.partitionReplicaAssignment(tp) - val liveReplicas = assignment.filter(replica => controllerContext.isReplicaOnline(replica, tp)) + val liveReplicas = assignment.filter(replica => controllerContextSnapshot.isReplicaOnline(replica, tp)) val isr = controllerContext.partitionLeadershipInfo(tp).get.leaderAndIsr.isr PartitionLeaderElectionAlgorithms .preferredReplicaPartitionLeaderElection(assignment, isr, liveReplicas.toSet) diff --git a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala index 20a6a87bd2b52..eb527b1937d38 100755 --- a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala +++ b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala @@ -81,12 +81,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(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 @@ -268,10 +269,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 } @@ -457,9 +459,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 f56d28dade458..d6f2076a8b3cc 100644 --- a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala +++ b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala @@ -56,11 +56,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. diff --git a/core/src/main/scala/kafka/controller/TopicDeletionManager.scala b/core/src/main/scala/kafka/controller/TopicDeletionManager.scala index eb60b1ce1b5ec..dc2e5ccce2e1d 100755 --- a/core/src/main/scala/kafka/controller/TopicDeletionManager.scala +++ b/core/src/main/scala/kafka/controller/TopicDeletionManager.scala @@ -308,10 +308,11 @@ class TopicDeletionManager(config: KafkaConfig, val allDeadReplicas = mutable.ListBuffer.empty[PartitionAndReplica] val allReplicasForDeletionRetry = mutable.ListBuffer.empty[PartitionAndReplica] val allTopicsIneligibleForDeletion = mutable.Set.empty[String] + val controllerContextSnapshot = ControllerContextSnapshot(controllerContext) topicsToBeDeleted.foreach { topic => val (aliveReplicas, deadReplicas) = controllerContext.replicasForTopic(topic).partition { r => - controllerContext.isReplicaOnline(r.replica, r.topicPartition) + controllerContextSnapshot.isReplicaOnline(r.replica, r.topicPartition) } val successfullyDeletedReplicas = controllerContext.replicasInState(topic, ReplicaDeletionSuccessful)