Repository: kafka
Updated Branches:
  refs/heads/trunk 9b8d99b1b -> 3f0168cd3


KAFKA-4229; Controller can't start after several zk expired event

Author: pengwei <pengwei.lihuawei.com>

Reviewers: wangguoz.gmail.com

Author: pengwei-li <[email protected]>
Author: c00353482 <[email protected]>

Reviewers: Jun Rao <[email protected]>

Closes #2175 from pengwei-li/trunk


Project: http://git-wip-us.apache.org/repos/asf/kafka/repo
Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/3f0168cd
Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/3f0168cd
Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/3f0168cd

Branch: refs/heads/trunk
Commit: 3f0168cd35c1baaa5a032e9a534f28a1bd9221b9
Parents: 9b8d99b
Author: pengwei-li <[email protected]>
Authored: Tue Jan 24 13:56:30 2017 -0800
Committer: Jun Rao <[email protected]>
Committed: Tue Jan 24 13:56:30 2017 -0800

----------------------------------------------------------------------
 .../main/scala/kafka/controller/KafkaController.scala    | 11 ++++++++---
 .../main/scala/kafka/server/ZookeeperLeaderElector.scala |  2 +-
 2 files changed, 9 insertions(+), 4 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/kafka/blob/3f0168cd/core/src/main/scala/kafka/controller/KafkaController.scala
----------------------------------------------------------------------
diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala 
b/core/src/main/scala/kafka/controller/KafkaController.scala
old mode 100755
new mode 100644
index 8ffa610..b0ed8d7
--- a/core/src/main/scala/kafka/controller/KafkaController.scala
+++ b/core/src/main/scala/kafka/controller/KafkaController.scala
@@ -1163,9 +1163,14 @@ class KafkaController(val config: KafkaConfig, zkUtils: 
ZkUtils, val brokerState
     @throws[Exception]
     def handleNewSession() {
       info("ZK expired; shut down all controller components and try to 
re-elect")
-      onControllerResignation()
-      inLock(controllerContext.controllerLock) {
-        controllerElector.elect
+      if (controllerElector.getControllerID() != config.brokerId) {
+        onControllerResignation()
+        inLock(controllerContext.controllerLock) {
+          controllerElector.elect
+        }
+      } else {
+        //maybe create by current session or the previous zk session's 
ephemeral node is not deleted
+        info("ZK expired, but the current controller id %d is the same as this 
broker id, skip re-elect".format(config.brokerId))
       }
     }
 

http://git-wip-us.apache.org/repos/asf/kafka/blob/3f0168cd/core/src/main/scala/kafka/server/ZookeeperLeaderElector.scala
----------------------------------------------------------------------
diff --git a/core/src/main/scala/kafka/server/ZookeeperLeaderElector.scala 
b/core/src/main/scala/kafka/server/ZookeeperLeaderElector.scala
old mode 100755
new mode 100644
index ca0f6a0..f41782e
--- a/core/src/main/scala/kafka/server/ZookeeperLeaderElector.scala
+++ b/core/src/main/scala/kafka/server/ZookeeperLeaderElector.scala
@@ -52,7 +52,7 @@ class ZookeeperLeaderElector(controllerContext: 
ControllerContext,
     }
   }
 
-  private def getControllerID(): Int = {
+  def getControllerID(): Int = {
     controllerContext.zkUtils.readDataMaybeNull(electionPath)._1 match {
        case Some(controller) => KafkaController.parseControllerId(controller)
        case None => -1

Reply via email to