From dd9a697d46f8b4578e364e72c30f5bfb67379ef6 Mon Sep 17 00:00:00 2001 From: Ke Hu Date: Wed, 27 Mar 2019 16:02:49 -0700 Subject: [PATCH 1/2] Set callback to null in addStopReplicaRequestForBrokers when replica state changes to offline --- .../src/main/scala/kafka/controller/ReplicaStateMachine.scala | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala index 433ab5668379e..58e40dca85e50 100644 --- a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala +++ b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala @@ -202,9 +202,11 @@ class ReplicaStateMachine(config: KafkaConfig, replicaState.put(replica, OnlineReplica) } case OfflineReplica => + // Set callback to null so that controller will send one grouped request instead of one request for + // each topic partition validReplicas.foreach { replica => controllerBrokerRequestBatch.addStopReplicaRequestForBrokers(Seq(replicaId), replica.topicPartition, - deletePartition = false, (_, _) => ()) + deletePartition = false, null) } val (replicasWithLeadershipInfo, replicasWithoutLeadershipInfo) = validReplicas.partition { replica => controllerContext.partitionLeadershipInfo.contains(replica.topicPartition) From 0bb43e561e4d6afbc2b279449be91a86f41ab8d3 Mon Sep 17 00:00:00 2001 From: Ke Hu Date: Thu, 28 Mar 2019 16:38:49 -0700 Subject: [PATCH 2/2] Set callback to null in addStopReplicaRequestForBrokers when replica state changes to offline --- .../scala/kafka/controller/ControllerChannelManager.scala | 2 +- .../src/main/scala/kafka/controller/ReplicaStateMachine.scala | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala index 3776b6903041a..c92d3d9558930 100755 --- a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala @@ -357,7 +357,7 @@ class ControllerBrokerRequestBatch(controller: KafkaController, stateChangeLogge stopReplicaRequestMap.getOrElseUpdate(brokerId, Seq.empty[StopReplicaRequestInfo]) val v = stopReplicaRequestMap(brokerId) stopReplicaRequestMap(brokerId) = v :+ StopReplicaRequestInfo(PartitionAndReplica(topicPartition, brokerId), - deletePartition, (r: AbstractResponse) => callback(r, brokerId)) + deletePartition, if (callback != null) (r: AbstractResponse) => callback(r, brokerId) else null) } } diff --git a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala index 433ab5668379e..58e40dca85e50 100644 --- a/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala +++ b/core/src/main/scala/kafka/controller/ReplicaStateMachine.scala @@ -202,9 +202,11 @@ class ReplicaStateMachine(config: KafkaConfig, replicaState.put(replica, OnlineReplica) } case OfflineReplica => + // Set callback to null so that controller will send one grouped request instead of one request for + // each topic partition validReplicas.foreach { replica => controllerBrokerRequestBatch.addStopReplicaRequestForBrokers(Seq(replicaId), replica.topicPartition, - deletePartition = false, (_, _) => ()) + deletePartition = false, null) } val (replicasWithLeadershipInfo, replicasWithoutLeadershipInfo) = validReplicas.partition { replica => controllerContext.partitionLeadershipInfo.contains(replica.topicPartition)