Intro and AI Disclaimer
I am a long time storm user of a small storm cluster. In general, the cluster stays robust and stable over the course of weeks but every now and then we run into a situation where seemingly an error in a topology and, following up, a zookeeper connection loss brings a supervisor down and then the whole cluster becomes unstable as topologies get re-assigned, supervisors blacklisted, and we have to wipe the whole thing and start fresh.
Since I never really could pin that down to a specific reason I used local AI to dig through our logs and the storm code. It came up with the report seen below, it also suggested a fix which I am happy to open a PR for, if that is okay in this repo. We currently have the fix running in production, seems to hold up so far.
Depending on the assessment of the storm dev team I can also overwork the fix and implement a more suitable solution, should the proposed new class ConnectionAwareRetryPolicy in storm-client not be ideal.
Everything below is written by AI but carefully reviewed by me.
Bug Description
A transient ZooKeeper connection loss (e.g. a single ensemble member closing its socket) causes the entire supervisor process to exit within ~600ms. All topology workers on that node are killed as a result.
Storm version: 3.1.0
Component: storm-server (supervisor), storm-client (CuratorUtils / ZK retry)
Expected Behavior
The supervisor should survive transient ZK connection blips. The ZK client's SendThread is already failovering to another ensemble member - the supervisor should wait for that reconnection and retry the operation, not die.
Actual Behavior
The supervisor process exits with code 20 within 629ms of the connection loss.
Root Cause
Three compounding issues:
-
Curator RetryLoop races the ZK SendThread. Both operate on the same TCP connection. The RetryLoop (driven by StormBoundedExponentialBackoffRetry) blindly sleeps between retries without checking whether the SendThread has already reconnected to another ensemble member.
-
The retry window is too short. With defaults (storm.zookeeper.retry.times=5, storm.zookeeper.retry.interval=1000), the total retry budget is ~30s. In practice the exception escapes before even one retry sleep completes, because the RetryLoop and SendThread are racing on the same connection.
-
DefaultUncaughtExceptionHandler kills the process. Any uncaught Throwable on a StormTimer thread calls Runtime.getRuntime().exit(20). A transient KeeperException$ConnectionLoss is treated the same as a fatal bug.
Evidence
From a supervisor daemon log:
08:39:53.907 ZK server closes socket (EndOfStreamException)
08:39:54.489 Curator ConnectionState → SUSPENDED (+582ms)
08:39:54.510 SupervisorHeartbeat → existsNode() → KeeperException$ConnectionLoss (+603ms)
08:39:54.536 DefaultUncaughtExceptionHandler → Utils.exitProcess(20) (+629ms) ← JVM exits
The ZK ensemble was healthy throughout - the restarted supervisor reconnected to the same ensemble within 100ms.
Proposed Fix
Two layers:
1. ConnectionAwareRetryPolicy (new class in storm-client)
A RetryPolicy wrapper that tracks the ConnectionState via a ConnectionStateListener. When the connection is SUSPENDED or LOST, it calls CuratorFramework.blockUntilConnected(sessionTimeout) to yield to the SendThread's failover instead of blind sleep+retry. On reconnection, the operation retries immediately on the new connection. Uses storm.zookeeper.session.timeout as the wait bound.
Wired up in CuratorUtils.newCurator() via AtomicReference late binding.
2. Safety net in SupervisorHeartbeat
Wrap the heartbeat body in a try/catch so that a total ZK outage (beyond the session timeout) skips the cycle instead of killing the process. A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 30s) allows 6 missed beats.
Files Changed
storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java (new)
storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
storm-server/src/main/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeat.java
Intro and AI Disclaimer
I am a long time storm user of a small storm cluster. In general, the cluster stays robust and stable over the course of weeks but every now and then we run into a situation where seemingly an error in a topology and, following up, a zookeeper connection loss brings a supervisor down and then the whole cluster becomes unstable as topologies get re-assigned, supervisors blacklisted, and we have to wipe the whole thing and start fresh.
Since I never really could pin that down to a specific reason I used local AI to dig through our logs and the storm code. It came up with the report seen below, it also suggested a fix which I am happy to open a PR for, if that is okay in this repo. We currently have the fix running in production, seems to hold up so far.
Depending on the assessment of the storm dev team I can also overwork the fix and implement a more suitable solution, should the proposed new class
ConnectionAwareRetryPolicyin storm-client not be ideal.Everything below is written by AI but carefully reviewed by me.
Bug Description
A transient ZooKeeper connection loss (e.g. a single ensemble member closing its socket) causes the entire supervisor process to exit within ~600ms. All topology workers on that node are killed as a result.
Storm version: 3.1.0
Component:
storm-server(supervisor),storm-client(CuratorUtils / ZK retry)Expected Behavior
The supervisor should survive transient ZK connection blips. The ZK client's
SendThreadis already failovering to another ensemble member - the supervisor should wait for that reconnection and retry the operation, not die.Actual Behavior
The supervisor process exits with code 20 within 629ms of the connection loss.
Root Cause
Three compounding issues:
Curator
RetryLoopraces the ZKSendThread. Both operate on the same TCP connection. TheRetryLoop(driven byStormBoundedExponentialBackoffRetry) blindly sleeps between retries without checking whether theSendThreadhas already reconnected to another ensemble member.The retry window is too short. With defaults (
storm.zookeeper.retry.times=5,storm.zookeeper.retry.interval=1000), the total retry budget is ~30s. In practice the exception escapes before even one retry sleep completes, because theRetryLoopandSendThreadare racing on the same connection.DefaultUncaughtExceptionHandlerkills the process. Any uncaughtThrowableon aStormTimerthread callsRuntime.getRuntime().exit(20). A transientKeeperException$ConnectionLossis treated the same as a fatal bug.Evidence
From a supervisor daemon log:
The ZK ensemble was healthy throughout - the restarted supervisor reconnected to the same ensemble within 100ms.
Proposed Fix
Two layers:
1.
ConnectionAwareRetryPolicy(new class instorm-client)A
RetryPolicywrapper that tracks theConnectionStatevia aConnectionStateListener. When the connection isSUSPENDEDorLOST, it callsCuratorFramework.blockUntilConnected(sessionTimeout)to yield to theSendThread's failover instead of blind sleep+retry. On reconnection, the operation retries immediately on the new connection. Usesstorm.zookeeper.session.timeoutas the wait bound.Wired up in
CuratorUtils.newCurator()viaAtomicReferencelate binding.2. Safety net in
SupervisorHeartbeatWrap the heartbeat body in a
try/catchso that a total ZK outage (beyond the session timeout) skips the cycle instead of killing the process. A missed heartbeat is harmless:nimbus.supervisor.timeout.secs(default 30s) allows 6 missed beats.Files Changed
storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java(new)storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.javastorm-server/src/main/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeat.java