diff --git a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala index 8228795a3d7c3..f1403238cbb35 100755 --- a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala @@ -400,7 +400,7 @@ abstract class AbstractControllerBrokerRequestBatch(config: KafkaConfig, val leaderEpoch = if (controllerContext.isTopicQueuedUpForDeletion(topicPartition.topic)) { LeaderAndIsr.EpochDuringDelete } else { - controllerContext.partitionLeadershipInfo.get(topicPartition) + controllerContext.partitionLeadershipInfo(topicPartition) .map(_.leaderAndIsr.leaderEpoch) .getOrElse(LeaderAndIsr.NoEpoch) } @@ -420,7 +420,7 @@ abstract class AbstractControllerBrokerRequestBatch(config: KafkaConfig, partitions: collection.Set[TopicPartition]): Unit = { def updateMetadataRequestPartitionInfo(partition: TopicPartition, beingDeleted: Boolean): Unit = { - controllerContext.partitionLeadershipInfo.get(partition) match { + controllerContext.partitionLeadershipInfo(partition) match { case Some(LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) => val replicas = controllerContext.partitionReplicaAssignment(partition) val offlineReplicas = replicas.filter(!controllerContext.isReplicaOnline(_, partition)) diff --git a/core/src/main/scala/kafka/controller/ControllerContext.scala b/core/src/main/scala/kafka/controller/ControllerContext.scala index 51986a4abc6a5..530df2f3e18e9 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -74,6 +74,7 @@ case class ReplicaAssignment private (replicas: Seq[Int], class ControllerContext { val stats = new ControllerStats var offlinePartitionCount = 0 + var preferredReplicaImbalanceCount = 0 val shuttingDownBrokerIds = mutable.Set.empty[Int] private val liveBrokers = mutable.Set.empty[Broker] private val liveBrokerEpochs = mutable.Map.empty[Int, Long] @@ -82,7 +83,7 @@ class ControllerContext { val allTopics = mutable.Set.empty[String] val partitionAssignments = mutable.Map.empty[String, mutable.Map[Int, ReplicaAssignment]] - val partitionLeadershipInfo = mutable.Map.empty[TopicPartition, LeaderIsrAndControllerEpoch] + private val partitionLeadershipInfo = mutable.Map.empty[TopicPartition, LeaderIsrAndControllerEpoch] val partitionsBeingReassigned = mutable.Set.empty[TopicPartition] val partitionStates = mutable.Map.empty[TopicPartition, PartitionState] val replicaStates = mutable.Map.empty[PartitionAndReplica, ReplicaState] @@ -120,6 +121,7 @@ class ControllerContext { replicasOnOfflineDirs.clear() partitionStates.clear() offlinePartitionCount = 0 + preferredReplicaImbalanceCount = 0 replicaStates.clear() } @@ -137,7 +139,10 @@ class ControllerContext { def updatePartitionFullReplicaAssignment(topicPartition: TopicPartition, newAssignment: ReplicaAssignment): Unit = { val assignments = partitionAssignments.getOrElseUpdate(topicPartition.topic, mutable.Map.empty) - assignments.put(topicPartition.partition, newAssignment) + val previous = assignments.put(topicPartition.partition, newAssignment) + val leadershipInfo = partitionLeadershipInfo.get(topicPartition) + updatePreferredReplicaImbalanceMetric(topicPartition, previous, leadershipInfo, + Some(newAssignment), leadershipInfo) } def partitionReplicaAssignmentForTopic(topic : String): Map[TopicPartition, Seq[Int]] = { @@ -281,16 +286,24 @@ class ControllerContext { } def removeTopic(topic: String): Unit = { + // Metric is cleaned when the topic is queued up for deletion so + // we don't clean it twice. We clean it only if it is deleted + // directly. + if (!topicsToBeDeleted.contains(topic)) + cleanPreferredReplicaImbalanceMetric(topic) + topicsToBeDeleted -= topic + topicsWithDeletionStarted -= topic allTopics -= topic - partitionAssignments.remove(topic) - partitionLeadershipInfo.foreach { case (topicPartition, _) => - if (topicPartition.topic == topic) - partitionLeadershipInfo.remove(topicPartition) + partitionAssignments.remove(topic).foreach { assignments => + assignments.keys.foreach { partition => + partitionLeadershipInfo.remove(new TopicPartition(topic, partition)) + } } } def queueTopicDeletion(topics: Set[String]): Unit = { topicsToBeDeleted ++= topics + topics.foreach(cleanPreferredReplicaImbalanceMetric) } def beginTopicDeletion(topics: Set[String]): Unit = { @@ -391,6 +404,90 @@ class ControllerContext { partitionsForTopic(topic).filter { partition => states.contains(partitionState(partition)) }.toSet } + def putPartitionLeadershipInfo(partition: TopicPartition, + leaderIsrAndControllerEpoch: LeaderIsrAndControllerEpoch): Unit = { + val previous = partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + val replicaAssignment = partitionFullReplicaAssignment(partition) + updatePreferredReplicaImbalanceMetric(partition, Some(replicaAssignment), previous, + Some(replicaAssignment), Some(leaderIsrAndControllerEpoch)) + } + + def partitionLeadershipInfo(partition: TopicPartition): Option[LeaderIsrAndControllerEpoch] = { + partitionLeadershipInfo.get(partition) + } + + def partitionsLeadershipInfo(): Iterable[(TopicPartition, LeaderIsrAndControllerEpoch)] = { + partitionLeadershipInfo + } + + def partitionsWithLeaders(): Set[TopicPartition] = { + partitionLeadershipInfo.keys.filter(tp => !isTopicQueuedUpForDeletion(tp.topic)).toSet + } + + def partitionsWithOfflineLeader(): Set[TopicPartition] = { + partitionLeadershipInfo.filter { case (topicPartition, leaderIsrAndControllerEpoch) => + !isReplicaOnline(leaderIsrAndControllerEpoch.leaderAndIsr.leader, topicPartition) && + !isTopicQueuedUpForDeletion(topicPartition.topic) + }.keySet + } + + def partitionLeadersOnBroker(brokerId: Int): Set[TopicPartition] = { + partitionLeadershipInfo.filter { case (topicPartition, leaderIsrAndControllerEpoch) => + !isTopicQueuedUpForDeletion(topicPartition.topic) && + leaderIsrAndControllerEpoch.leaderAndIsr.leader == brokerId && + partitionReplicaAssignment(topicPartition).size > 1 + }.keySet + } + + def clearPartitionLeadershipInfo(): Unit = { + partitionLeadershipInfo.clear() + } + + def partitionWithLeadersCount(): Int = { + partitionLeadershipInfo.size + } + + private def updatePreferredReplicaImbalanceMetric(partition: TopicPartition, + oldReplicaAssignment: Option[ReplicaAssignment], + oldLeadershipInfo: Option[LeaderIsrAndControllerEpoch], + newReplicaAssignment: Option[ReplicaAssignment], + newLeadershipInfo: Option[LeaderIsrAndControllerEpoch]): Unit = { + if (!isTopicQueuedUpForDeletion(partition.topic)) { + oldReplicaAssignment.foreach { replicaAssignment => + oldLeadershipInfo.foreach { leadershipInfo => + if (!hasPreferredLeader(replicaAssignment, leadershipInfo)) + preferredReplicaImbalanceCount -= 1 + } + } + + newReplicaAssignment.foreach { replicaAssignment => + newLeadershipInfo.foreach { leadershipInfo => + if (!hasPreferredLeader(replicaAssignment, leadershipInfo)) + preferredReplicaImbalanceCount += 1 + } + } + } + } + + private def cleanPreferredReplicaImbalanceMetric(topic: String): Unit = { + partitionAssignments.getOrElse(topic, mutable.Map.empty).foreach { case (partition, replicaAssignment) => + partitionLeadershipInfo.get(new TopicPartition(topic, partition)).foreach { leadershipInfo => + if (!hasPreferredLeader(replicaAssignment, leadershipInfo)) + preferredReplicaImbalanceCount -= 1 + } + } + } + + private def hasPreferredLeader(replicaAssignment: ReplicaAssignment, + leadershipInfo: LeaderIsrAndControllerEpoch): Boolean = { + val preferredReplica = replicaAssignment.replicas.head + if (replicaAssignment.isBeingReassigned && replicaAssignment.addingReplicas.contains(preferredReplica)) + // reassigning partitions are not counted as imbalanced until the new replica joins the ISR (completes reassignment) + !leadershipInfo.leaderAndIsr.isr.contains(preferredReplica) + else + leadershipInfo.leaderAndIsr.leader == preferredReplica + } + private def isValidReplicaStateTransition(replica: PartitionAndReplica, targetState: ReplicaState): Boolean = targetState.validPreviousStates.contains(replicaStates(replica)) diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index ab2b600ad4758..f5757b37d5502 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -354,7 +354,7 @@ class KafkaController(val config: KafkaConfig, // Send update metadata request to all the new brokers in the cluster with a full set of partition states for initialization. // In cases of controlled shutdown leaders will not be elected when a new broker comes up. So at least in the // common controlled shutdown case, the metadata will reach the new brokers faster. - sendUpdateMetadataRequest(newBrokers, controllerContext.partitionLeadershipInfo.keySet) + sendUpdateMetadataRequest(newBrokers, controllerContext.partitionsWithLeaders) // the very first thing to do when a new broker comes up is send it the entire list of partitions that it is // supposed to host. Based on that the broker starts the high watermark threads for the input list of partitions val allReplicasOnNewBrokers = controllerContext.replicasOnBrokers(newBrokersSet) @@ -439,12 +439,10 @@ class KafkaController(val config: KafkaConfig, val (newOfflineReplicasForDeletion, newOfflineReplicasNotForDeletion) = newOfflineReplicas.partition(p => topicDeletionManager.isTopicQueuedUpForDeletion(p.topic)) - val partitionsWithoutLeader = controllerContext.partitionLeadershipInfo.filter(partitionAndLeader => - !controllerContext.isReplicaOnline(partitionAndLeader._2.leaderAndIsr.leader, partitionAndLeader._1) && - !topicDeletionManager.isTopicQueuedUpForDeletion(partitionAndLeader._1.topic)).keySet + val partitionsWithOfflineLeader = controllerContext.partitionsWithOfflineLeader() // trigger OfflinePartition state for all partitions whose current leader is one amongst the newOfflineReplicas - partitionStateMachine.handleStateChanges(partitionsWithoutLeader.toSeq, OfflinePartition) + partitionStateMachine.handleStateChanges(partitionsWithOfflineLeader.toSeq, OfflinePartition) // trigger OnlinePartition state changes for offline or new partitions partitionStateMachine.triggerOnlinePartitionStateChange() // trigger OfflineReplica state change for those newly offline replicas @@ -460,7 +458,7 @@ class KafkaController(val config: KafkaConfig, // If replica failure did not require leader re-election, inform brokers of the offline brokers // Note that during leader re-election, brokers update their metadata - if (partitionsWithoutLeader.isEmpty) { + if (partitionsWithOfflineLeader.isEmpty) { sendUpdateMetadataRequest(controllerContext.liveOrShuttingDownBrokerIds.toSeq, Set.empty) } } @@ -734,7 +732,7 @@ class KafkaController(val config: KafkaConfig, if (replicaAssignment.isBeingReassigned) controllerContext.partitionsBeingReassigned.add(topicPartition) } - controllerContext.partitionLeadershipInfo.clear() + controllerContext.clearPartitionLeadershipInfo() controllerContext.shuttingDownBrokerIds.clear() // register broker modifications handlers registerBrokerModificationsHandler(controllerContext.liveOrShuttingDownBrokerIds) @@ -754,7 +752,7 @@ class KafkaController(val config: KafkaConfig, val replicas = controllerContext.partitionReplicaAssignment(partition) val topicDeleted = replicas.isEmpty val successful = - if (!topicDeleted) controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader == replicas.head else false + if (!topicDeleted) controllerContext.partitionLeadershipInfo(partition).get.leaderAndIsr.leader == replicas.head else false successful || topicDeleted } val pendingPreferredReplicaElectionsIgnoringTopicDeletion = partitionsUndergoingPreferredReplicaElection -- partitionsThatCompletedPreferredReplicaElection @@ -796,7 +794,7 @@ class KafkaController(val config: KafkaConfig, private def updateLeaderAndIsrCache(partitions: Seq[TopicPartition] = controllerContext.allPartitions.toSeq): Unit = { val leaderIsrAndControllerEpochs = zkClient.getTopicPartitionStates(partitions) leaderIsrAndControllerEpochs.foreach { case (partition, leaderIsrAndControllerEpoch) => - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) } } @@ -815,7 +813,7 @@ class KafkaController(val config: KafkaConfig, private def moveReassignedPartitionLeaderIfRequired(topicPartition: TopicPartition, newAssignment: ReplicaAssignment): Unit = { val reassignedReplicas = newAssignment.replicas - val currentLeader = controllerContext.partitionLeadershipInfo(topicPartition).leaderAndIsr.leader + val currentLeader = controllerContext.partitionLeadershipInfo(topicPartition).get.leaderAndIsr.leader if (!reassignedReplicas.contains(currentLeader)) { info(s"Leader $currentLeader for partition $topicPartition being reassigned, " + @@ -961,7 +959,7 @@ class KafkaController(val config: KafkaConfig, isTriggeredByAutoRebalance : Boolean): Unit = { for (partition <- partitionsToBeRemoved) { // check the status - val currentLeader = controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader + val currentLeader = controllerContext.partitionLeadershipInfo(partition).get.leaderAndIsr.leader val preferredReplica = controllerContext.partitionReplicaAssignment(partition).head if (currentLeader == preferredReplica) { info(s"Partition $partition completed preferred replica leader election. New leader is $preferredReplica") @@ -1024,7 +1022,7 @@ class KafkaController(val config: KafkaConfig, finishedUpdates.get(partition) match { case Some(Right(leaderAndIsr)) => val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, epoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) finalLeaderIsrAndControllerEpoch = Some(leaderIsrAndControllerEpoch) info(s"Updated leader epoch for partition $partition to ${leaderAndIsr.leaderEpoch}") true @@ -1051,7 +1049,7 @@ class KafkaController(val config: KafkaConfig, // for each broker, check if a preferred replica election needs to be triggered preferredReplicasForTopicsByBrokers.foreach { case (leaderBroker, topicPartitionsForBroker) => val topicsNotInPreferredReplica = topicPartitionsForBroker.filter { case (topicPartition, _) => - val leadershipInfo = controllerContext.partitionLeadershipInfo.get(topicPartition) + val leadershipInfo = controllerContext.partitionLeadershipInfo(topicPartition) leadershipInfo.exists(_.leaderAndIsr.leader != leaderBroker) } debug(s"Topics not in preferred replica for broker $leaderBroker $topicsNotInPreferredReplica") @@ -1078,7 +1076,7 @@ class KafkaController(val config: KafkaConfig, 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 + val isr = controllerContext.partitionLeadershipInfo(tp).get.leaderAndIsr.isr PartitionLeaderElectionAlgorithms .preferredReplicaPartitionLeaderElection(assignment, isr, liveReplicas.toSet) .nonEmpty @@ -1143,11 +1141,11 @@ class KafkaController(val config: KafkaConfig, val partitionsToActOn = controllerContext.partitionsOnBroker(id).filter { partition => controllerContext.partitionReplicaAssignment(partition).size > 1 && - controllerContext.partitionLeadershipInfo.contains(partition) && + controllerContext.partitionLeadershipInfo(partition).isDefined && !topicDeletionManager.isTopicQueuedUpForDeletion(partition.topic) } val (partitionsLedByBroker, partitionsFollowedByBroker) = partitionsToActOn.partition { partition => - controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader == id + controllerContext.partitionLeadershipInfo(partition).get.leaderAndIsr.leader == id } partitionStateMachine.handleStateChanges(partitionsLedByBroker.toSeq, OnlinePartition, Some(ControlledShutdownPartitionLeaderElectionStrategy)) try { @@ -1163,16 +1161,8 @@ class KafkaController(val config: KafkaConfig, // If the broker is a follower, updates the isr in ZK and notifies the current leader replicaStateMachine.handleStateChanges(partitionsFollowedByBroker.map(partition => PartitionAndReplica(partition, id)).toSeq, OfflineReplica) - def replicatedPartitionsBrokerLeads() = { - trace(s"All leaders = ${controllerContext.partitionLeadershipInfo.mkString(",")}") - controllerContext.partitionLeadershipInfo.filter { - case (topicPartition, leaderIsrAndControllerEpoch) => - !topicDeletionManager.isTopicQueuedUpForDeletion(topicPartition.topic) && - leaderIsrAndControllerEpoch.leaderAndIsr.leader == id && - controllerContext.partitionReplicaAssignment(topicPartition).size > 1 - }.keys - } - replicatedPartitionsBrokerLeads().toSet + trace(s"All leaders = ${controllerContext.partitionsLeadershipInfo.mkString(",")}") + controllerContext.partitionLeadersOnBroker(id) } private def processUpdateMetadataResponseReceived(updateMetadataResponse: UpdateMetadataResponse, brokerId: Int): Unit = { @@ -1254,28 +1244,12 @@ class KafkaController(val config: KafkaConfig, if (!isActive) { 0 } else { - controllerContext.allPartitions.count { topicPartition => - val replicaAssignment: ReplicaAssignment = controllerContext.partitionFullReplicaAssignment(topicPartition) - val replicas = replicaAssignment.replicas - val preferredReplica = replicas.head - - val isImbalanced = controllerContext.partitionLeadershipInfo.get(topicPartition) match { - case Some(leadershipInfo) => - if (replicaAssignment.isBeingReassigned && replicaAssignment.addingReplicas.contains(preferredReplica)) - // reassigning partitions are not counted as imbalanced until the new replica joins the ISR (completes reassignment) - leadershipInfo.leaderAndIsr.isr.contains(preferredReplica) - else - leadershipInfo.leaderAndIsr.leader != preferredReplica - case None => false - } - - isImbalanced && !topicDeletionManager.isTopicQueuedUpForDeletion(topicPartition.topic) - } + controllerContext.preferredReplicaImbalanceCount } globalTopicCount = if (!isActive) 0 else controllerContext.allTopics.size - globalPartitionCount = if (!isActive) 0 else controllerContext.partitionLeadershipInfo.size + globalPartitionCount = if (!isActive) 0 else controllerContext.partitionWithLeadersCount topicsToDeleteCount = if (!isActive) 0 else controllerContext.topicsToBeDeleted.size @@ -1760,11 +1734,11 @@ class KafkaController(val config: KafkaConfig, case ElectionType.PREFERRED => val assignedReplicas = controllerContext.partitionReplicaAssignment(partition) val preferredReplica = assignedReplicas.head - val currentLeader = controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader + val currentLeader = controllerContext.partitionLeadershipInfo(partition).get.leaderAndIsr.leader currentLeader != preferredReplica case ElectionType.UNCLEAN => - val currentLeader = controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader + val currentLeader = controllerContext.partitionLeadershipInfo(partition).get.leaderAndIsr.leader currentLeader == LeaderAndIsr.NoLeader || !controllerContext.liveBrokerIds.contains(currentLeader) } } diff --git a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala index eba2478e6588f..4bd2579346a0a 100755 --- a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala +++ b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala @@ -82,7 +82,7 @@ abstract class PartitionStateMachine(controllerContext: ControllerContext) exten private def initializePartitionState(): Unit = { 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 { + 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)) @@ -226,7 +226,7 @@ class ZkPartitionStateMachine(config: KafkaConfig, val successfulInitializations = initializeLeaderAndIsrForPartitions(uninitializedPartitions) successfulInitializations.foreach { partition => stateChangeLog.info(s"Changed partition $partition from ${partitionState(partition)} to $targetState with state " + - s"${controllerContext.partitionLeadershipInfo(partition).leaderAndIsr}") + s"${controllerContext.partitionLeadershipInfo(partition).get.leaderAndIsr}") controllerContext.putPartitionState(partition, OnlinePartition) } } @@ -309,7 +309,7 @@ class ZkPartitionStateMachine(config: KafkaConfig, val partition = createResponse.ctx.get.asInstanceOf[TopicPartition] val leaderIsrAndControllerEpoch = leaderIsrAndControllerEpochs(partition) if (code == Code.OK) { - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) controllerBrokerRequestBatch.addLeaderAndIsrRequestForBrokers(leaderIsrAndControllerEpoch.leaderAndIsr.isr, partition, leaderIsrAndControllerEpoch, controllerContext.partitionFullReplicaAssignment(partition), isNew = true) successfulInitializations += partition @@ -437,7 +437,7 @@ class ZkPartitionStateMachine(config: KafkaConfig, result.foreach { leaderAndIsr => val replicaAssignment = controllerContext.partitionFullReplicaAssignment(partition) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerContext.epoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) controllerBrokerRequestBatch.addLeaderAndIsrRequestForBrokers(recipientsPerPartition(partition), partition, leaderIsrAndControllerEpoch, replicaAssignment, isNew = false) } diff --git a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala index a0df78ba3710c..3fbc4f42fef8c 100644 --- a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala +++ b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala @@ -168,7 +168,7 @@ class ZkReplicaStateMachine(config: KafkaConfig, val partition = replica.topicPartition val currentState = controllerContext.replicaState(replica) - controllerContext.partitionLeadershipInfo.get(partition) match { + controllerContext.partitionLeadershipInfo(partition) match { case Some(leaderIsrAndControllerEpoch) => if (leaderIsrAndControllerEpoch.leaderAndIsr.leader == replicaId) { val exception = new StateChangeFailedException(s"Replica $replicaId for partition $partition cannot be moved to NewReplica state as it is being requested to become leader") @@ -203,7 +203,7 @@ class ZkReplicaStateMachine(config: KafkaConfig, controllerContext.updatePartitionFullReplicaAssignment(partition, newAssignment) } case _ => - controllerContext.partitionLeadershipInfo.get(partition) match { + controllerContext.partitionLeadershipInfo(partition) match { case Some(leaderIsrAndControllerEpoch) => controllerBrokerRequestBatch.addLeaderAndIsrRequestForBrokers(Seq(replicaId), replica.topicPartition, @@ -221,7 +221,7 @@ class ZkReplicaStateMachine(config: KafkaConfig, controllerBrokerRequestBatch.addStopReplicaRequestForBrokers(Seq(replicaId), replica.topicPartition, deletePartition = false) } val (replicasWithLeadershipInfo, replicasWithoutLeadershipInfo) = validReplicas.partition { replica => - controllerContext.partitionLeadershipInfo.contains(replica.topicPartition) + controllerContext.partitionLeadershipInfo(replica.topicPartition).isDefined } val updatedLeaderIsrAndControllerEpochs = removeReplicasFromIsr(replicaId, replicasWithLeadershipInfo.map(_.topicPartition)) updatedLeaderIsrAndControllerEpochs.foreach { case (partition, leaderIsrAndControllerEpoch) => @@ -364,7 +364,7 @@ class ZkReplicaStateMachine(config: KafkaConfig, (leaderAndIsrsWithoutReplica ++ finishedPartitions).map { case (partition, result) => (partition, result.map { leaderAndIsr => val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerContext.epoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) leaderIsrAndControllerEpoch }) } diff --git a/core/src/main/scala/kafka/controller/TopicDeletionManager.scala b/core/src/main/scala/kafka/controller/TopicDeletionManager.scala index a9639b46a3445..0fd7274d128ad 100755 --- a/core/src/main/scala/kafka/controller/TopicDeletionManager.scala +++ b/core/src/main/scala/kafka/controller/TopicDeletionManager.scala @@ -243,8 +243,6 @@ class TopicDeletionManager(config: KafkaConfig, val replicasForDeletedTopic = controllerContext.replicasInState(topic, ReplicaDeletionSuccessful) // controller will remove this replica from the state machine as well as its partition assignment cache replicaStateMachine.handleStateChanges(replicasForDeletedTopic.toSeq, NonExistentReplica) - controllerContext.topicsToBeDeleted -= topic - controllerContext.topicsWithDeletionStarted -= topic client.deleteTopic(topic, controllerContext.epochZkVersion) controllerContext.removeTopic(topic) } diff --git a/core/src/test/scala/unit/kafka/controller/ControllerChannelManagerTest.scala b/core/src/test/scala/unit/kafka/controller/ControllerChannelManagerTest.scala index 60802862260da..58c7203794196 100644 --- a/core/src/test/scala/unit/kafka/controller/ControllerChannelManagerTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ControllerChannelManagerTest.scala @@ -61,7 +61,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, leaderAndIsr) => val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - context.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + context.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) batch.addLeaderAndIsrRequestForBrokers(Seq(2), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(1, 2, 3)), isNew = false) } batch.sendRequestsToBrokers(controllerEpoch) @@ -99,7 +99,7 @@ class ControllerChannelManagerTest { val leaderAndIsr = LeaderAndIsr(1, List(1, 2)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - context.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + context.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) batch.newBatch() batch.addLeaderAndIsrRequestForBrokers(Seq(2), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(1, 2, 3)), isNew = true) @@ -131,7 +131,7 @@ class ControllerChannelManagerTest { val leaderAndIsr = LeaderAndIsr(1, List(1, 2)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - context.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + context.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) batch.newBatch() batch.addLeaderAndIsrRequestForBrokers(Seq(1, 2, 3), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(1, 2, 3)), isNew = false) @@ -177,7 +177,7 @@ class ControllerChannelManagerTest { val leaderAndIsr = LeaderAndIsr(1, List(1, 2)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - context.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + context.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) batch.newBatch() batch.addLeaderAndIsrRequestForBrokers(Seq(2), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(1, 2, 3)), isNew = false) @@ -201,7 +201,7 @@ class ControllerChannelManagerTest { ) partitions.foreach { case (partition, leaderAndIsr) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) } batch.newBatch() @@ -272,7 +272,7 @@ class ControllerChannelManagerTest { ) partitions.foreach { case (partition, leaderAndIsr) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) } context.queueTopicDeletion(Set("foo")) @@ -380,7 +380,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -409,7 +409,7 @@ class ControllerChannelManagerTest { val partition = new TopicPartition("foo", 0) val leaderAndIsr = LeaderAndIsr(1, List(1, 2)) - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.newBatch() batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition = true) @@ -447,7 +447,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -494,7 +494,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -549,7 +549,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -608,11 +608,11 @@ class ControllerChannelManagerTest { batch.newBatch() deletePartitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition) } nonDeletePartitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -662,7 +662,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2, 3), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -704,7 +704,7 @@ class ControllerChannelManagerTest { batch.newBatch() partitions.foreach { case (partition, LeaderAndDelete(leaderAndIsr, deletePartition)) => - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.addStopReplicaRequestForBrokers(Seq(2, 3), partition, deletePartition) } batch.sendRequestsToBrokers(controllerEpoch) @@ -748,7 +748,7 @@ class ControllerChannelManagerTest { val partition = new TopicPartition("foo", 0) val leaderAndIsr = LeaderAndIsr(1, List(1, 2)) - context.partitionLeadershipInfo.put(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) + context.putPartitionLeadershipInfo(partition, LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) batch.newBatch() batch.addStopReplicaRequestForBrokers(Seq(2), partition, deletePartition = false) diff --git a/core/src/test/scala/unit/kafka/controller/ControllerContextTest.scala b/core/src/test/scala/unit/kafka/controller/ControllerContextTest.scala index 96240b25ea5e4..dbc91e97c1273 100644 --- a/core/src/test/scala/unit/kafka/controller/ControllerContextTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ControllerContextTest.scala @@ -17,7 +17,9 @@ package unit.kafka.controller +import kafka.api.LeaderAndIsr import kafka.cluster.{Broker, EndPoint} +import kafka.controller.LeaderIsrAndControllerEpoch import kafka.controller.{ControllerContext, ReplicaAssignment} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.network.ListenerName @@ -160,4 +162,35 @@ class ControllerContextTest { assertEquals(assignment, firstReassign.reassignTo(Seq(1,2,3))) } + @Test + def testPreferredReplicaImbalanceMetric(): Unit = { + context.updatePartitionFullReplicaAssignment(tp1, ReplicaAssignment(Seq(1, 2, 3))) + context.updatePartitionFullReplicaAssignment(tp2, ReplicaAssignment(Seq(1, 2, 3))) + context.updatePartitionFullReplicaAssignment(tp3, ReplicaAssignment(Seq(1, 2, 3))) + assertEquals(0, context.preferredReplicaImbalanceCount) + + context.putPartitionLeadershipInfo(tp1, LeaderIsrAndControllerEpoch(LeaderAndIsr(1, List(1, 2, 3)), 0)) + assertEquals(0, context.preferredReplicaImbalanceCount) + + context.putPartitionLeadershipInfo(tp2, LeaderIsrAndControllerEpoch(LeaderAndIsr(2, List(2, 3, 1)), 0)) + assertEquals(1, context.preferredReplicaImbalanceCount) + + context.putPartitionLeadershipInfo(tp3, LeaderIsrAndControllerEpoch(LeaderAndIsr(3, List(3, 1, 2)), 0)) + assertEquals(2, context.preferredReplicaImbalanceCount) + + context.updatePartitionFullReplicaAssignment(tp1, ReplicaAssignment(Seq(2, 3, 1))) + context.updatePartitionFullReplicaAssignment(tp2, ReplicaAssignment(Seq(2, 3, 1))) + assertEquals(2, context.preferredReplicaImbalanceCount) + + context.queueTopicDeletion(Set(tp3.topic)) + assertEquals(1, context.preferredReplicaImbalanceCount) + + context.putPartitionLeadershipInfo(tp3, LeaderIsrAndControllerEpoch(LeaderAndIsr(1, List(3, 1, 2)), 0)) + assertEquals(1, context.preferredReplicaImbalanceCount) + + context.removeTopic(tp1.topic) + context.removeTopic(tp2.topic) + context.removeTopic(tp3.topic) + assertEquals(0, context.preferredReplicaImbalanceCount) + } } diff --git a/core/src/test/scala/unit/kafka/controller/MockPartitionStateMachine.scala b/core/src/test/scala/unit/kafka/controller/MockPartitionStateMachine.scala index b29a3d9a585a0..b9a4d04198da0 100644 --- a/core/src/test/scala/unit/kafka/controller/MockPartitionStateMachine.scala +++ b/core/src/test/scala/unit/kafka/controller/MockPartitionStateMachine.scala @@ -85,7 +85,7 @@ class MockPartitionStateMachine(controllerContext: ControllerContext, val validLeaderAndIsrs = mutable.Buffer.empty[(TopicPartition, LeaderAndIsr)] for (partition <- partitions) { - val leaderIsrAndControllerEpoch = controllerContext.partitionLeadershipInfo(partition) + val leaderIsrAndControllerEpoch = controllerContext.partitionLeadershipInfo(partition).get if (leaderIsrAndControllerEpoch.controllerEpoch > controllerContext.epoch) { val failMsg = s"Aborted leader election for partition $partition since the LeaderAndIsr path was " + s"already written by another controller. This probably means that the current controller went through " + @@ -118,7 +118,7 @@ class MockPartitionStateMachine(controllerContext: ControllerContext, Left(new StateChangeFailedException(failMsg)) case Some(leaderAndIsr) => val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerContext.epoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) Right(leaderAndIsr) } diff --git a/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala b/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala index ae5bf48e92248..894d56e8dd7f9 100644 --- a/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala +++ b/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala @@ -162,7 +162,7 @@ class PartitionStateMachineTest { controllerContext.putPartitionState(partition, OnlinePartition) val leaderAndIsr = LeaderAndIsr(brokerId, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) val stat = new Stat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -198,7 +198,7 @@ class PartitionStateMachineTest { controllerContext.putPartitionState(partition, OnlinePartition) val leaderAndIsr = LeaderAndIsr(brokerId, List(brokerId, otherBrokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) val stat = new Stat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -251,7 +251,7 @@ class PartitionStateMachineTest { controllerContext.putPartitionState(partition, OfflinePartition) val leaderAndIsr = LeaderAndIsr(LeaderAndIsr.NoLeader, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) val stat = new Stat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -298,7 +298,7 @@ class PartitionStateMachineTest { val leaderAndIsr = LeaderAndIsr(leaderBrokerId, List(leaderBrokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) EasyMock @@ -356,7 +356,7 @@ class PartitionStateMachineTest { controllerContext.putPartitionState(partition, OfflinePartition) val leaderAndIsr = LeaderAndIsr(LeaderAndIsr.NoLeader, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) EasyMock.expect(mockZkClient.getTopicPartitionStatesRaw(partitions)) @@ -381,7 +381,7 @@ class PartitionStateMachineTest { controllerContext.putPartitionState(partition, OfflinePartition) val leaderAndIsr = LeaderAndIsr(LeaderAndIsr.NoLeader, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) val stat = new Stat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) diff --git a/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala b/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala index e79d2f581e5aa..cbf6df3da86b6 100644 --- a/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala @@ -191,7 +191,7 @@ class ReplicaStateMachineTest { controllerContext.putReplicaState(replica, OnlineReplica) controllerContext.updatePartitionFullReplicaAssignment(partition, ReplicaAssignment(Seq(brokerId))) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) EasyMock.expect(mockControllerBrokerRequestBatch.addLeaderAndIsrRequestForBrokers(Seq(brokerId), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(brokerId)), isNew = false)) @@ -210,7 +210,7 @@ class ReplicaStateMachineTest { controllerContext.updatePartitionFullReplicaAssignment(partition, ReplicaAssignment(replicaIds)) val leaderAndIsr = LeaderAndIsr(brokerId, replicaIds) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) val stat = new Stat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -230,7 +230,7 @@ class ReplicaStateMachineTest { EasyMock.replay(mockZkClient, mockControllerBrokerRequestBatch) replicaStateMachine.handleStateChanges(replicas, OfflineReplica) EasyMock.verify(mockZkClient, mockControllerBrokerRequestBatch) - assertEquals(updatedLeaderIsrAndControllerEpoch, controllerContext.partitionLeadershipInfo(partition)) + assertEquals(updatedLeaderIsrAndControllerEpoch, controllerContext.partitionLeadershipInfo(partition).get) assertEquals(OfflineReplica, replicaState(replica)) } @@ -264,7 +264,7 @@ class ReplicaStateMachineTest { controllerContext.putReplicaState(replica, OfflineReplica) controllerContext.updatePartitionFullReplicaAssignment(partition, ReplicaAssignment(Seq(brokerId))) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) EasyMock.expect(mockControllerBrokerRequestBatch.addLeaderAndIsrRequestForBrokers(Seq(brokerId), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(brokerId)), isNew = false)) @@ -382,7 +382,7 @@ class ReplicaStateMachineTest { controllerContext.putReplicaState(replica, ReplicaDeletionIneligible) controllerContext.updatePartitionFullReplicaAssignment(partition, ReplicaAssignment(Seq(brokerId))) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) - controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + controllerContext.putPartitionLeadershipInfo(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) EasyMock.expect(mockControllerBrokerRequestBatch.addLeaderAndIsrRequestForBrokers(Seq(brokerId), partition, leaderIsrAndControllerEpoch, replicaAssignment(Seq(brokerId)), isNew = false))