Repository: kafka Updated Branches: refs/heads/0.10.2 cd3040767 -> 3a169837b
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 (cherry picked from commit 3f0168cd35c1baaa5a032e9a534f28a1bd9221b9) Signed-off-by: Jun Rao <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/kafka/repo Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/3a169837 Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/3a169837 Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/3a169837 Branch: refs/heads/0.10.2 Commit: 3a169837b864886f19c57dfaa6c9b3a188ffbb4d Parents: cd30407 Author: pengwei-li <[email protected]> Authored: Tue Jan 24 13:56:30 2017 -0800 Committer: Jun Rao <[email protected]> Committed: Tue Jan 24 13:57:58 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/3a169837/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/3a169837/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
