From b2177a645798a5c5699b088f590bd15dcbed6b6d Mon Sep 17 00:00:00 2001 From: Lucas Wang Date: Thu, 8 Mar 2018 13:59:07 -0800 Subject: [PATCH 1/3] Speeds up the processing of StopReplicaResponse events on the controller --- .../kafka/controller/ControllerContext.scala | 68 ++++++++++--- .../kafka/controller/KafkaController.scala | 98 ++++++++++--------- .../controller/PartitionStateMachine.scala | 2 +- .../controller/ReplicaStateMachine.scala | 7 +- .../controller/TopicDeletionManager.scala | 3 +- .../PartitionStateMachineTest.scala | 16 +-- .../controller/ReplicaStateMachineTest.scala | 12 +-- 7 files changed, 128 insertions(+), 78 deletions(-) diff --git a/core/src/main/scala/kafka/controller/ControllerContext.scala b/core/src/main/scala/kafka/controller/ControllerContext.scala index 541bce82d4870..dcc63555d6460 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -31,7 +31,7 @@ class ControllerContext { var epoch: Int = KafkaController.InitialControllerEpoch - 1 var epochZkVersion: Int = KafkaController.InitialControllerEpochZkVersion - 1 var allTopics: Set[String] = Set.empty - var partitionReplicaAssignment: mutable.Map[TopicPartition, Seq[Int]] = mutable.Map.empty + private var partitionReplicaAssignmentUnderlying: mutable.Map[String, mutable.Map[Int, Seq[Int]]] = mutable.Map.empty var partitionLeadershipInfo: mutable.Map[TopicPartition, LeaderIsrAndControllerEpoch] = mutable.Map.empty val partitionsBeingReassigned: mutable.Map[TopicPartition, ReassignedPartitionsContext] = mutable.Map.empty val replicasOnOfflineDirs: mutable.Map[Int, Set[TopicPartition]] = mutable.Map.empty @@ -39,6 +39,38 @@ class ControllerContext { private var liveBrokersUnderlying: Set[Broker] = Set.empty private var liveBrokerIdsUnderlying: Set[Int] = Set.empty + def partitionReplicaAssignment(topicPartition: TopicPartition) : Seq[Int] = { + partitionReplicaAssignmentUnderlying.getOrElse(topicPartition.topic, mutable.Map.empty) + .getOrElse(topicPartition.partition, Seq.empty) + } + + def clearPartitionReplicaAssignment() = { + partitionReplicaAssignmentUnderlying = mutable.Map.empty + } + + def updatePartitionReplicaAssignment(topicPartition: TopicPartition, newReplicas : Seq[Int]) = { + partitionReplicaAssignmentUnderlying.getOrElseUpdate(topicPartition.topic, mutable.Map.empty) + .put(topicPartition.partition, newReplicas) + } + + def partitionReplicaAssignmentForTopic(topic : String) : mutable.Map[TopicPartition, Seq[Int]] = { + partitionReplicaAssignmentUnderlying.getOrElse(topic, mutable.Map.empty).map { + case (partition, replicas) => (new TopicPartition(topic, partition), replicas) + } + } + + def removePartitionReplicaAssignmentForTopic(topic : String) = { + partitionReplicaAssignmentUnderlying.remove(topic) + } + + def allPartitions : Set[TopicPartition] = { + partitionReplicaAssignmentUnderlying.flatMap { + case (topic, topicReplicaAssignment) => topicReplicaAssignment.map { + case (partition, _) => new TopicPartition(topic, partition) + } + }.toSet + } + // setter def liveBrokers_=(brokers: Set[Broker]) { liveBrokersUnderlying = brokers @@ -53,8 +85,12 @@ class ControllerContext { def liveOrShuttingDownBrokers = liveBrokersUnderlying def partitionsOnBroker(brokerId: Int): Set[TopicPartition] = { - partitionReplicaAssignment.collect { - case (topicPartition, replicas) if replicas.contains(brokerId) => topicPartition + partitionReplicaAssignmentUnderlying.flatMap { + case (topic, topicReplicaAssignment) => topicReplicaAssignment.filter { + case (_, replicas) => replicas.contains(brokerId) + }.map { + case (partition, _) => new TopicPartition(topic, partition) + } }.toSet } @@ -68,22 +104,26 @@ class ControllerContext { def replicasOnBrokers(brokerIds: Set[Int]): Set[PartitionAndReplica] = { brokerIds.flatMap { brokerId => - partitionReplicaAssignment.collect { case (topicPartition, replicas) if replicas.contains(brokerId) => - PartitionAndReplica(topicPartition, brokerId) + partitionReplicaAssignmentUnderlying.flatMap { + case (topic, topicReplicaAssignment) => topicReplicaAssignment.collect { + case (partition, replicas) if replicas.contains(brokerId) => + PartitionAndReplica(new TopicPartition(topic, partition), brokerId) + } } - }.toSet + } } def replicasForTopic(topic: String): Set[PartitionAndReplica] = { - partitionReplicaAssignment - .filter { case (topicPartition, _) => topicPartition.topic == topic } - .flatMap { case (topicPartition, replicas) => - replicas.map(PartitionAndReplica(topicPartition, _)) - }.toSet + partitionReplicaAssignmentUnderlying.getOrElse(topic, mutable.Map.empty).flatMap { + case (partition, replicas) => replicas.map(r => PartitionAndReplica(new TopicPartition(topic, partition), r)) + }.toSet } - def partitionsForTopic(topic: String): collection.Set[TopicPartition] = - partitionReplicaAssignment.keySet.filter(topicPartition => topicPartition.topic == topic) + def partitionsForTopic(topic: String): collection.Set[TopicPartition] = { + partitionReplicaAssignmentUnderlying.getOrElse(topic, mutable.Map.empty).map { + case (partition, _) => new TopicPartition(topic, partition) + }.toSet + } def allLiveReplicas(): Set[PartitionAndReplica] = { replicasOnBrokers(liveBrokerIds).filter { partitionAndReplica => @@ -100,7 +140,7 @@ class ControllerContext { def removeTopic(topic: String) = { partitionLeadershipInfo = partitionLeadershipInfo.filter { case (topicPartition, _) => topicPartition.topic != topic } - partitionReplicaAssignment = partitionReplicaAssignment.filter { case (topicPartition, _) => topicPartition.topic != topic } + partitionReplicaAssignmentUnderlying.remove(topic) allTopics -= topic } diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index a8707ad887d76..daaef0d34f53d 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -569,28 +569,28 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti } val newReplicas = reassignedPartitionContext.newReplicas val topic = tp.topic - controllerContext.partitionReplicaAssignment.get(tp) match { - case Some(assignedReplicas) => - if (assignedReplicas == newReplicas) { - info(s"Partition $tp to be reassigned is already assigned to replicas " + - s"${newReplicas.mkString(",")}. Ignoring request for partition reassignment.") - removePartitionFromReassignedPartitions(tp) - } else { - try { - info(s"Handling reassignment of partition $tp to new replicas ${newReplicas.mkString(",")}") - // first register ISR change listener - reassignedPartitionContext.registerReassignIsrChangeHandler(zkClient) - // mark topic ineligible for deletion for the partitions being reassigned - topicDeletionManager.markTopicIneligibleForDeletion(Set(topic)) - onPartitionReassignment(tp, reassignedPartitionContext) - } catch { - case e: Throwable => - error(s"Error completing reassignment of partition $tp", e) - // remove the partition from the admin path to unblock the admin client - removePartitionFromReassignedPartitions(tp) - } + val assignedReplicas = controllerContext.partitionReplicaAssignment(tp) + if (assignedReplicas.nonEmpty) { + if (assignedReplicas == newReplicas) { + info(s"Partition $tp to be reassigned is already assigned to replicas " + + s"${newReplicas.mkString(",")}. Ignoring request for partition reassignment.") + removePartitionFromReassignedPartitions(tp) + } else { + try { + info(s"Handling reassignment of partition $tp to new replicas ${newReplicas.mkString(",")}") + // first register ISR change listener + reassignedPartitionContext.registerReassignIsrChangeHandler(zkClient) + // mark topic ineligible for deletion for the partitions being reassigned + topicDeletionManager.markTopicIneligibleForDeletion(Set(topic)) + onPartitionReassignment(tp, reassignedPartitionContext) + } catch { + case e: Throwable => + error(s"Error completing reassignment of partition $tp", e) + // remove the partition from the admin path to unblock the admin client + removePartitionFromReassignedPartitions(tp) } - case None => + } + } else { error(s"Ignoring request to reassign partition $tp that doesn't exist.") removePartitionFromReassignedPartitions(tp) } @@ -643,7 +643,9 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti controllerContext.liveBrokers = zkClient.getAllBrokersInCluster.toSet controllerContext.allTopics = zkClient.getAllTopicsInCluster.toSet registerPartitionModificationsHandlers(controllerContext.allTopics.toSeq) - controllerContext.partitionReplicaAssignment = mutable.Map.empty ++ zkClient.getReplicaAssignmentForTopics(controllerContext.allTopics.toSet) + zkClient.getReplicaAssignmentForTopics(controllerContext.allTopics.toSet).foreach { + case (topicPartition, assignedReplicas) => controllerContext.updatePartitionReplicaAssignment(topicPartition, assignedReplicas) + } controllerContext.partitionLeadershipInfo = new mutable.HashMap[TopicPartition, LeaderIsrAndControllerEpoch] controllerContext.shuttingDownBrokerIds = mutable.Set.empty[Int] // register broker modifications handlers @@ -662,10 +664,10 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti val partitionsUndergoingPreferredReplicaElection = zkClient.getPreferredReplicaElection // check if they are already completed or topic was deleted val partitionsThatCompletedPreferredReplicaElection = partitionsUndergoingPreferredReplicaElection.filter { partition => - val replicasOpt = controllerContext.partitionReplicaAssignment.get(partition) - val topicDeleted = replicasOpt.isEmpty + val replicas = controllerContext.partitionReplicaAssignment(partition) + val topicDeleted = replicas.isEmpty val successful = - if (!topicDeleted) controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader == replicasOpt.get.head else false + if (!topicDeleted) controllerContext.partitionLeadershipInfo(partition).leaderAndIsr.leader == replicas.head else false successful || topicDeleted } val pendingPreferredReplicaElectionsIgnoringTopicDeletion = partitionsUndergoingPreferredReplicaElection -- partitionsThatCompletedPreferredReplicaElection @@ -687,7 +689,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti controllerContext.epoch = 0 controllerContext.epochZkVersion = 0 controllerContext.allTopics = Set.empty - controllerContext.partitionReplicaAssignment.clear() + controllerContext.clearPartitionReplicaAssignment() controllerContext.partitionLeadershipInfo.clear() controllerContext.partitionsBeingReassigned.clear() controllerContext.liveBrokers = Set.empty @@ -706,9 +708,10 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti private def fetchTopicDeletionsInProgress(): (Set[String], Set[String]) = { val topicsToBeDeleted = zkClient.getTopicDeletions.toSet - val topicsWithOfflineReplicas = controllerContext.partitionReplicaAssignment.filter { case (partition, replicas) => - replicas.exists(r => !controllerContext.isReplicaOnline(r, partition)) - }.keySet.map(_.topic) + val topicsWithOfflineReplicas = controllerContext.allTopics.filter { topic => { + val replicasForTopic = controllerContext.replicasForTopic(topic) + replicasForTopic.exists(r => !controllerContext.isReplicaOnline(r.replica, new TopicPartition(topic, r.partition))) + }} val topicsForWhichPartitionReassignmentIsInProgress = controllerContext.partitionsBeingReassigned.keySet.map(_.topic) val topicsIneligibleForDeletion = topicsWithOfflineReplicas | topicsForWhichPartitionReassignmentIsInProgress info(s"List of topics to be deleted: ${topicsToBeDeleted.mkString(",")}") @@ -722,7 +725,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti controllerContext.controllerChannelManager.startup() } - private def updateLeaderAndIsrCache(partitions: Seq[TopicPartition] = controllerContext.partitionReplicaAssignment.keys.toSeq) { + private def updateLeaderAndIsrCache(partitions: Seq[TopicPartition] = controllerContext.allPartitions.toSeq) { val leaderIsrAndControllerEpochs = zkClient.getTopicPartitionStates(partitions) leaderIsrAndControllerEpochs.foreach { case (partition, leaderIsrAndControllerEpoch) => controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) @@ -742,7 +745,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti // change the assigned replica list to just the reassigned replicas in the cache so it gets sent out on the LeaderAndIsr // request to the current or new leader. This will prevent it from adding the old replicas to the ISR val oldAndNewReplicas = controllerContext.partitionReplicaAssignment(topicPartition) - controllerContext.partitionReplicaAssignment.put(topicPartition, reassignedReplicas) + controllerContext.updatePartitionReplicaAssignment(topicPartition, reassignedReplicas) if (!reassignedPartitionContext.newReplicas.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") @@ -778,14 +781,14 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti private def updateAssignedReplicasForPartition(partition: TopicPartition, replicas: Seq[Int]) { - val partitionsAndReplicasForThisTopic = controllerContext.partitionReplicaAssignment.filter(_._1.topic == partition.topic) + val partitionsAndReplicasForThisTopic = controllerContext.partitionReplicaAssignmentForTopic(partition.topic) partitionsAndReplicasForThisTopic.put(partition, replicas) val setDataResponse = zkClient.setTopicAssignmentRaw(partition.topic, partitionsAndReplicasForThisTopic.toMap) setDataResponse.resultCode match { case Code.OK => info(s"Updated assigned replicas for partition $partition being reassigned to ${replicas.mkString(",")}") // update the assigned replica list after a successful zookeeper write - controllerContext.partitionReplicaAssignment.put(partition, replicas) + controllerContext.updatePartitionReplicaAssignment(partition, replicas) case Code.NONODE => throw new IllegalStateException(s"Topic ${partition.topic} doesn't exist") case _ => throw new KafkaException(setDataResponse.resultException.get) } @@ -971,9 +974,12 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti private def checkAndTriggerAutoLeaderRebalance(): Unit = { trace("Checking need to trigger auto leader balancing") val preferredReplicasForTopicsByBrokers: Map[Int, Map[TopicPartition, Seq[Int]]] = - controllerContext.partitionReplicaAssignment.filterNot { case (tp, _) => - topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) - }.groupBy { case (_, assignedReplicas) => assignedReplicas.head } + controllerContext.allPartitions.filterNot { + tp => topicDeletionManager.isTopicQueuedUpForDeletion(tp.topic) + }.map { tp => + (tp, controllerContext.partitionReplicaAssignment(tp) ) + }.toMap.groupBy { case (_, assignedReplicas) => assignedReplicas.head } + debug(s"Preferred replicas by broker $preferredReplicasForTopicsByBrokers") // for each broker, check if a preferred replica election needs to be triggered @@ -1155,11 +1161,12 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti if (!isActive) { 0 } else { - controllerContext.partitionReplicaAssignment.count { case (topicPartition, replicas) => + controllerContext.allPartitions.count { topicAndPartition => + val replicas = controllerContext.partitionReplicaAssignment(topicAndPartition) val preferredReplica = replicas.head - val leadershipInfo = controllerContext.partitionLeadershipInfo.get(topicPartition) + val leadershipInfo = controllerContext.partitionLeadershipInfo.get(topicAndPartition) leadershipInfo.map(_.leaderAndIsr.leader != preferredReplica).getOrElse(false) && - !topicDeletionManager.isTopicQueuedUpForDeletion(topicPartition.topic) + !topicDeletionManager.isTopicQueuedUpForDeletion(topicAndPartition.topic) } } @@ -1279,9 +1286,10 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti registerPartitionModificationsHandlers(newTopics.toSeq) val addedPartitionReplicaAssignment = zkClient.getReplicaAssignmentForTopics(newTopics) - controllerContext.partitionReplicaAssignment = controllerContext.partitionReplicaAssignment.filter(p => - !deletedTopics.contains(p._1.topic)) - controllerContext.partitionReplicaAssignment ++= addedPartitionReplicaAssignment + deletedTopics.foreach(controllerContext.removePartitionReplicaAssignmentForTopic) + addedPartitionReplicaAssignment.foreach { + case (topicAndPartition, newReplicas) => controllerContext.updatePartitionReplicaAssignment(topicAndPartition, newReplicas) + } info(s"New topics: [$newTopics], deleted topics: [$deletedTopics], new partition replica assignment " + s"[$addedPartitionReplicaAssignment]") if (addedPartitionReplicaAssignment.nonEmpty) @@ -1312,14 +1320,16 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti if (!isActive) return val partitionReplicaAssignment = zkClient.getReplicaAssignmentForTopics(immutable.Set(topic)) val partitionsToBeAdded = partitionReplicaAssignment.filter(p => - !controllerContext.partitionReplicaAssignment.contains(p._1)) + controllerContext.partitionReplicaAssignment(p._1).isEmpty) if (topicDeletionManager.isTopicQueuedUpForDeletion(topic)) error(s"Skipping adding partitions ${partitionsToBeAdded.map(_._1.partition).mkString(",")} for topic $topic " + "since it is currently being deleted") else { if (partitionsToBeAdded.nonEmpty) { info(s"New partitions to be added $partitionsToBeAdded") - controllerContext.partitionReplicaAssignment ++= partitionsToBeAdded + partitionsToBeAdded.foreach { case (topicAndPartition, assignedReplicas) => + controllerContext.updatePartitionReplicaAssignment(topicAndPartition, assignedReplicas) + } onNewPartitionCreation(partitionsToBeAdded.keySet) } } diff --git a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala index 2e27272158f4c..74bc59faee2a5 100755 --- a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala +++ b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala @@ -76,7 +76,7 @@ class PartitionStateMachine(config: KafkaConfig, * zookeeper */ private def initializePartitionState() { - for (topicPartition <- controllerContext.partitionReplicaAssignment.keys) { + 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) => diff --git a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala index 85af764b635e2..a2d04e65ae6bc 100644 --- a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala +++ b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala @@ -80,7 +80,8 @@ class ReplicaStateMachine(config: KafkaConfig, * in zookeeper */ private def initializeReplicaState() { - controllerContext.partitionReplicaAssignment.foreach { case (partition, replicas) => + controllerContext.allPartitions.foreach { partition => + val replicas = controllerContext.partitionReplicaAssignment(partition) replicas.foreach { replicaId => val partitionAndReplica = PartitionAndReplica(partition, replicaId) if (controllerContext.isReplicaOnline(replicaId, partition)) @@ -181,7 +182,7 @@ class ReplicaStateMachine(config: KafkaConfig, case NewReplica => val assignment = controllerContext.partitionReplicaAssignment(partition) if (!assignment.contains(replicaId)) { - controllerContext.partitionReplicaAssignment.put(partition, assignment :+ replicaId) + controllerContext.updatePartitionReplicaAssignment(partition, assignment :+ replicaId) } case _ => controllerContext.partitionLeadershipInfo.get(partition) match { @@ -237,7 +238,7 @@ class ReplicaStateMachine(config: KafkaConfig, case NonExistentReplica => validReplicas.foreach { replica => val currentAssignedReplicas = controllerContext.partitionReplicaAssignment(replica.topicPartition) - controllerContext.partitionReplicaAssignment.put(replica.topicPartition, currentAssignedReplicas.filterNot(_ == replica.replica)) + controllerContext.updatePartitionReplicaAssignment(replica.topicPartition, currentAssignedReplicas.filterNot(_ == replica.replica)) logSuccessfulTransition(replicaId, replica.topicPartition, replicaState(replica), NonExistentReplica) replicaState.remove(replica) } diff --git a/core/src/main/scala/kafka/controller/TopicDeletionManager.scala b/core/src/main/scala/kafka/controller/TopicDeletionManager.scala index b1d83947b14ca..6e145516cce8d 100755 --- a/core/src/main/scala/kafka/controller/TopicDeletionManager.scala +++ b/core/src/main/scala/kafka/controller/TopicDeletionManager.scala @@ -255,9 +255,8 @@ class TopicDeletionManager(controller: KafkaController, // send update metadata so that brokers stop serving data for topics to be deleted val partitions = topics.flatMap(controllerContext.partitionsForTopic) controller.sendUpdateMetadataRequest(controllerContext.liveOrShuttingDownBrokerIds.toSeq, partitions) - val partitionReplicaAssignmentByTopic = controllerContext.partitionReplicaAssignment.groupBy(p => p._1.topic) topics.foreach { topic => - onPartitionDeletion(partitionReplicaAssignmentByTopic(topic).keySet) + onPartitionDeletion(controllerContext.partitionsForTopic(topic)) } } diff --git a/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala b/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala index 32e0d43084d28..52f459970d166 100644 --- a/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala +++ b/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala @@ -80,7 +80,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testNewPartitionToOnlinePartitionTransition(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, NewPartition) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -98,7 +98,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testNewPartitionToOnlinePartitionTransitionZkUtilsExceptionFromCreateStates(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, NewPartition) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -114,7 +114,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testNewPartitionToOnlinePartitionTransitionErrorCodeFromCreateStates(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, NewPartition) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -144,7 +144,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testOnlinePartitionToOnlineTransition(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, OnlinePartition) val leaderAndIsr = LeaderAndIsr(brokerId, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) @@ -175,7 +175,7 @@ class PartitionStateMachineTest extends JUnitSuite { val otherBrokerId = brokerId + 1 controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0), TestUtils.createBroker(otherBrokerId, "host", 0)) controllerContext.shuttingDownBrokerIds.add(brokerId) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId, otherBrokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId, otherBrokerId)) partitionState.put(partition, OnlinePartition) val leaderAndIsr = LeaderAndIsr(brokerId, List(brokerId, otherBrokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) @@ -226,7 +226,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testOfflinePartitionToOnlinePartitionTransition(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, OfflinePartition) val leaderAndIsr = LeaderAndIsr(LeaderAndIsr.NoLeader, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) @@ -257,7 +257,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testOfflinePartitionToOnlinePartitionTransitionZkUtilsExceptionFromStateLookup(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, OfflinePartition) val leaderAndIsr = LeaderAndIsr(LeaderAndIsr.NoLeader, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) @@ -278,7 +278,7 @@ class PartitionStateMachineTest extends JUnitSuite { @Test def testOfflinePartitionToOnlinePartitionTransitionErrorCodeFromStateLookup(): Unit = { controllerContext.liveBrokers = Set(TestUtils.createBroker(brokerId, "host", 0)) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) partitionState.put(partition, OfflinePartition) val leaderAndIsr = LeaderAndIsr(LeaderAndIsr.NoLeader, List(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) diff --git a/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala b/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala index 4d38aac1659e3..6a961a53157d4 100644 --- a/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala +++ b/core/src/test/scala/unit/kafka/controller/ReplicaStateMachineTest.scala @@ -104,7 +104,7 @@ class ReplicaStateMachineTest extends JUnitSuite { @Test def testNewReplicaToOnlineReplicaTransition(): Unit = { replicaState.put(replica, NewReplica) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) replicaStateMachine.handleStateChanges(replicas, OnlineReplica) assertEquals(OnlineReplica, replicaState(replica)) } @@ -150,7 +150,7 @@ class ReplicaStateMachineTest extends JUnitSuite { @Test def testOnlineReplicaToOnlineReplicaTransition(): Unit = { replicaState.put(replica, OnlineReplica) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -168,7 +168,7 @@ class ReplicaStateMachineTest extends JUnitSuite { val otherBrokerId = brokerId + 1 val replicaIds = List(brokerId, otherBrokerId) replicaState.put(replica, OnlineReplica) - controllerContext.partitionReplicaAssignment.put(partition, replicaIds) + controllerContext.updatePartitionReplicaAssignment(partition, replicaIds) val leaderAndIsr = LeaderAndIsr(brokerId, replicaIds) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch) controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) @@ -225,7 +225,7 @@ class ReplicaStateMachineTest extends JUnitSuite { @Test def testOfflineReplicaToOnlineReplicaTransition(): Unit = { replicaState.put(replica, OfflineReplica) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) @@ -299,7 +299,7 @@ class ReplicaStateMachineTest extends JUnitSuite { @Test def testReplicaDeletionSuccessfulToNonexistentReplicaTransition(): Unit = { replicaState.put(replica, ReplicaDeletionSuccessful) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) replicaStateMachine.handleStateChanges(replicas, NonExistentReplica) assertEquals(Seq.empty, controllerContext.partitionReplicaAssignment(partition)) assertEquals(None, replicaState.get(replica)) @@ -343,7 +343,7 @@ class ReplicaStateMachineTest extends JUnitSuite { @Test def testReplicaDeletionIneligibleToOnlineReplicaTransition(): Unit = { replicaState.put(replica, ReplicaDeletionIneligible) - controllerContext.partitionReplicaAssignment.put(partition, Seq(brokerId)) + controllerContext.updatePartitionReplicaAssignment(partition, Seq(brokerId)) val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(LeaderAndIsr(brokerId, List(brokerId)), controllerEpoch) controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) EasyMock.expect(mockControllerBrokerRequestBatch.newBatch()) From c2dc88adbfd3509baf9bcb0231dc05fff748128e Mon Sep 17 00:00:00 2001 From: Lucas Wang Date: Fri, 23 Mar 2018 20:55:49 -0700 Subject: [PATCH 2/3] Addressed Jun's comments --- .../kafka/controller/ControllerContext.scala | 42 +++++++++++------ .../kafka/controller/KafkaController.scala | 45 +++++++------------ 2 files changed, 43 insertions(+), 44 deletions(-) diff --git a/core/src/main/scala/kafka/controller/ControllerContext.scala b/core/src/main/scala/kafka/controller/ControllerContext.scala index dcc63555d6460..373bb11199456 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -32,7 +32,7 @@ class ControllerContext { var epochZkVersion: Int = KafkaController.InitialControllerEpochZkVersion - 1 var allTopics: Set[String] = Set.empty private var partitionReplicaAssignmentUnderlying: mutable.Map[String, mutable.Map[Int, Seq[Int]]] = mutable.Map.empty - var partitionLeadershipInfo: mutable.Map[TopicPartition, LeaderIsrAndControllerEpoch] = mutable.Map.empty + val partitionLeadershipInfo: mutable.Map[TopicPartition, LeaderIsrAndControllerEpoch] = mutable.Map.empty val partitionsBeingReassigned: mutable.Map[TopicPartition, ReassignedPartitionsContext] = mutable.Map.empty val replicasOnOfflineDirs: mutable.Map[Int, Set[TopicPartition]] = mutable.Map.empty @@ -44,23 +44,23 @@ class ControllerContext { .getOrElse(topicPartition.partition, Seq.empty) } - def clearPartitionReplicaAssignment() = { - partitionReplicaAssignmentUnderlying = mutable.Map.empty + def clearTopicsState(): Unit = { + allTopics = Set.empty + partitionReplicaAssignmentUnderlying.clear() + partitionLeadershipInfo.clear() + partitionsBeingReassigned.clear() + replicasOnOfflineDirs.clear() } - def updatePartitionReplicaAssignment(topicPartition: TopicPartition, newReplicas : Seq[Int]) = { + def updatePartitionReplicaAssignment(topicPartition: TopicPartition, newReplicas: Seq[Int]) : Unit = { partitionReplicaAssignmentUnderlying.getOrElseUpdate(topicPartition.topic, mutable.Map.empty) .put(topicPartition.partition, newReplicas) } - def partitionReplicaAssignmentForTopic(topic : String) : mutable.Map[TopicPartition, Seq[Int]] = { - partitionReplicaAssignmentUnderlying.getOrElse(topic, mutable.Map.empty).map { + def partitionReplicaAssignmentForTopic(topic : String) : Map[TopicPartition, Seq[Int]] = { + partitionReplicaAssignmentUnderlying.getOrElse(topic, Map.empty).map { case (partition, replicas) => (new TopicPartition(topic, partition), replicas) - } - } - - def removePartitionReplicaAssignmentForTopic(topic : String) = { - partitionReplicaAssignmentUnderlying.remove(topic) + }.toMap } def allPartitions : Set[TopicPartition] = { @@ -138,10 +138,24 @@ class ControllerContext { } } + def resetContext() : Unit = { + if (controllerChannelManager != null) { + controllerChannelManager.shutdown() + controllerChannelManager = null + } + shuttingDownBrokerIds.clear() + epoch = 0 + epochZkVersion = 0 + clearTopicsState() + liveBrokers = Set.empty + } + def removeTopic(topic: String) = { - partitionLeadershipInfo = partitionLeadershipInfo.filter { case (topicPartition, _) => topicPartition.topic != topic } - partitionReplicaAssignmentUnderlying.remove(topic) allTopics -= topic + partitionReplicaAssignmentUnderlying.remove(topic) + partitionLeadershipInfo.foreach { + case (topicPartition, _) if topicPartition.topic == topic => partitionLeadershipInfo.remove(topicPartition) + case _ => + } } - } diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index daaef0d34f53d..ff0b9ec9bc2f8 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -309,7 +309,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti replicaStateMachine.shutdown() zkClient.unregisterZNodeChildChangeHandler(brokerChangeHandler.path) - resetControllerContext() + controllerContext.resetContext() info("Resigned") } @@ -646,7 +646,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti zkClient.getReplicaAssignmentForTopics(controllerContext.allTopics.toSet).foreach { case (topicPartition, assignedReplicas) => controllerContext.updatePartitionReplicaAssignment(topicPartition, assignedReplicas) } - controllerContext.partitionLeadershipInfo = new mutable.HashMap[TopicPartition, LeaderIsrAndControllerEpoch] + controllerContext.partitionLeadershipInfo.clear() controllerContext.shuttingDownBrokerIds = mutable.Set.empty[Int] // register broker modifications handlers registerBrokerModificationsHandler(controllerContext.liveBrokers.map(_.id)) @@ -680,21 +680,6 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti pendingPreferredReplicaElections } - private def resetControllerContext(): Unit = { - if (controllerContext.controllerChannelManager != null) { - controllerContext.controllerChannelManager.shutdown() - controllerContext.controllerChannelManager = null - } - controllerContext.shuttingDownBrokerIds.clear() - controllerContext.epoch = 0 - controllerContext.epochZkVersion = 0 - controllerContext.allTopics = Set.empty - controllerContext.clearPartitionReplicaAssignment() - controllerContext.partitionLeadershipInfo.clear() - controllerContext.partitionsBeingReassigned.clear() - controllerContext.liveBrokers = Set.empty - } - private def initializePartitionReassignment() { // read the partitions being reassigned from zookeeper path /admin/reassign_partitions val partitionsBeingReassigned = zkClient.getPartitionReassignment @@ -710,7 +695,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti val topicsToBeDeleted = zkClient.getTopicDeletions.toSet val topicsWithOfflineReplicas = controllerContext.allTopics.filter { topic => { val replicasForTopic = controllerContext.replicasForTopic(topic) - replicasForTopic.exists(r => !controllerContext.isReplicaOnline(r.replica, new TopicPartition(topic, r.partition))) + replicasForTopic.exists(r => !controllerContext.isReplicaOnline(r.replica, r.topicPartition)) }} val topicsForWhichPartitionReassignmentIsInProgress = controllerContext.partitionsBeingReassigned.keySet.map(_.topic) val topicsIneligibleForDeletion = topicsWithOfflineReplicas | topicsForWhichPartitionReassignmentIsInProgress @@ -781,9 +766,8 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti private def updateAssignedReplicasForPartition(partition: TopicPartition, replicas: Seq[Int]) { - val partitionsAndReplicasForThisTopic = controllerContext.partitionReplicaAssignmentForTopic(partition.topic) - partitionsAndReplicasForThisTopic.put(partition, replicas) - val setDataResponse = zkClient.setTopicAssignmentRaw(partition.topic, partitionsAndReplicasForThisTopic.toMap) + controllerContext.updatePartitionReplicaAssignment(partition, replicas) + val setDataResponse = zkClient.setTopicAssignmentRaw(partition.topic, controllerContext.partitionReplicaAssignmentForTopic(partition.topic)) setDataResponse.resultCode match { case Code.OK => info(s"Updated assigned replicas for partition $partition being reassigned to ${replicas.mkString(",")}") @@ -1161,12 +1145,12 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti if (!isActive) { 0 } else { - controllerContext.allPartitions.count { topicAndPartition => - val replicas = controllerContext.partitionReplicaAssignment(topicAndPartition) + controllerContext.allPartitions.count { topicPartition => + val replicas = controllerContext.partitionReplicaAssignment(topicPartition) val preferredReplica = replicas.head - val leadershipInfo = controllerContext.partitionLeadershipInfo.get(topicAndPartition) + val leadershipInfo = controllerContext.partitionLeadershipInfo.get(topicPartition) leadershipInfo.map(_.leaderAndIsr.leader != preferredReplica).getOrElse(false) && - !topicDeletionManager.isTopicQueuedUpForDeletion(topicAndPartition.topic) + !topicDeletionManager.isTopicQueuedUpForDeletion(topicPartition.topic) } } @@ -1286,7 +1270,7 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti registerPartitionModificationsHandlers(newTopics.toSeq) val addedPartitionReplicaAssignment = zkClient.getReplicaAssignmentForTopics(newTopics) - deletedTopics.foreach(controllerContext.removePartitionReplicaAssignmentForTopic) + deletedTopics.foreach(controllerContext.removeTopic) addedPartitionReplicaAssignment.foreach { case (topicAndPartition, newReplicas) => controllerContext.updatePartitionReplicaAssignment(topicAndPartition, newReplicas) } @@ -1319,16 +1303,17 @@ class KafkaController(val config: KafkaConfig, zkClient: KafkaZkClient, time: Ti override def process(): Unit = { if (!isActive) return val partitionReplicaAssignment = zkClient.getReplicaAssignmentForTopics(immutable.Set(topic)) - val partitionsToBeAdded = partitionReplicaAssignment.filter(p => - controllerContext.partitionReplicaAssignment(p._1).isEmpty) + val partitionsToBeAdded = partitionReplicaAssignment.filter { case (topicPartition, _) => + controllerContext.partitionReplicaAssignment(topicPartition).isEmpty + } if (topicDeletionManager.isTopicQueuedUpForDeletion(topic)) error(s"Skipping adding partitions ${partitionsToBeAdded.map(_._1.partition).mkString(",")} for topic $topic " + "since it is currently being deleted") else { if (partitionsToBeAdded.nonEmpty) { info(s"New partitions to be added $partitionsToBeAdded") - partitionsToBeAdded.foreach { case (topicAndPartition, assignedReplicas) => - controllerContext.updatePartitionReplicaAssignment(topicAndPartition, assignedReplicas) + partitionsToBeAdded.foreach { case (topicPartition, assignedReplicas) => + controllerContext.updatePartitionReplicaAssignment(topicPartition, assignedReplicas) } onNewPartitionCreation(partitionsToBeAdded.keySet) } From 28bc4ead5ed4d38fa5544e68473f141bf9be8b20 Mon Sep 17 00:00:00 2001 From: Lucas Wang Date: Thu, 29 Mar 2018 12:33:05 -0700 Subject: [PATCH 3/3] Making clearTopicsState private and removing whitespaces --- .../scala/kafka/controller/ControllerContext.scala | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/core/src/main/scala/kafka/controller/ControllerContext.scala b/core/src/main/scala/kafka/controller/ControllerContext.scala index 373bb11199456..f4671cfaa1a63 100644 --- a/core/src/main/scala/kafka/controller/ControllerContext.scala +++ b/core/src/main/scala/kafka/controller/ControllerContext.scala @@ -39,12 +39,12 @@ class ControllerContext { private var liveBrokersUnderlying: Set[Broker] = Set.empty private var liveBrokerIdsUnderlying: Set[Int] = Set.empty - def partitionReplicaAssignment(topicPartition: TopicPartition) : Seq[Int] = { + def partitionReplicaAssignment(topicPartition: TopicPartition): Seq[Int] = { partitionReplicaAssignmentUnderlying.getOrElse(topicPartition.topic, mutable.Map.empty) .getOrElse(topicPartition.partition, Seq.empty) } - def clearTopicsState(): Unit = { + private def clearTopicsState(): Unit = { allTopics = Set.empty partitionReplicaAssignmentUnderlying.clear() partitionLeadershipInfo.clear() @@ -52,18 +52,18 @@ class ControllerContext { replicasOnOfflineDirs.clear() } - def updatePartitionReplicaAssignment(topicPartition: TopicPartition, newReplicas: Seq[Int]) : Unit = { + def updatePartitionReplicaAssignment(topicPartition: TopicPartition, newReplicas: Seq[Int]): Unit = { partitionReplicaAssignmentUnderlying.getOrElseUpdate(topicPartition.topic, mutable.Map.empty) .put(topicPartition.partition, newReplicas) } - def partitionReplicaAssignmentForTopic(topic : String) : Map[TopicPartition, Seq[Int]] = { + def partitionReplicaAssignmentForTopic(topic : String): Map[TopicPartition, Seq[Int]] = { partitionReplicaAssignmentUnderlying.getOrElse(topic, Map.empty).map { case (partition, replicas) => (new TopicPartition(topic, partition), replicas) }.toMap } - def allPartitions : Set[TopicPartition] = { + def allPartitions: Set[TopicPartition] = { partitionReplicaAssignmentUnderlying.flatMap { case (topic, topicReplicaAssignment) => topicReplicaAssignment.map { case (partition, _) => new TopicPartition(topic, partition) @@ -138,7 +138,7 @@ class ControllerContext { } } - def resetContext() : Unit = { + def resetContext(): Unit = { if (controllerChannelManager != null) { controllerChannelManager.shutdown() controllerChannelManager = null @@ -150,7 +150,7 @@ class ControllerContext { liveBrokers = Set.empty } - def removeTopic(topic: String) = { + def removeTopic(topic: String): Unit = { allTopics -= topic partitionReplicaAssignmentUnderlying.remove(topic) partitionLeadershipInfo.foreach {