Conversation
Root cause: A transient ZooKeeper connection loss during the supervisor heartbeat cycle caused the process to die in 629ms via DefaultUncaughtExceptionHandler. The Curator RetryLoop blindly slept between retries, racing against the ZK client's SendThread reconnection instead of yielding to it. Fix (two layers): 1. ConnectionAwareRetryPolicy (new): A RetryPolicy wrapper that tracks the ZK 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 (no magic numbers). 2. SupervisorHeartbeat safety net: Wrap the heartbeat body in try/catch so that a total ZK outage (beyond session timeout) skips the cycle instead of killing the process. A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (30s) allows 6 missed beats. Wired up in CuratorUtils.newCurator() via AtomicReference late binding (the CuratorFramework doesn't exist until builder.build() returns).
reiabreu
left a comment
There was a problem hiding this comment.
This review was generated with the help of an LLM (Claude) and reviewed by me before posting.
Thanks for this — the #9098 diagnosis is solid and the two-layer approach (connection-aware retry + heartbeat safety net) targets the real root cause. I pulled the branch and ran the storm-client/storm-server suites. Findings below, with code where it helps. One blocker.
1. This breaks an existing test (blocker)
CuratorUtilsTest.newCuratorUsesExponentialBackoffTest casts the framework's retry policy directly to StormBoundedExponentialBackoffRetry. Since newCurator() now returns a ConnectionAwareRetryPolicy wrapper, that cast throws ClassCastException and mvn -pl storm-client test fails:
java.lang.ClassCastException: class ConnectionAwareRetryPolicy cannot be cast to class StormBoundedExponentialBackoffRetry at CuratorUtilsTest.java:75
The fix keeps the original coverage (that newCurator wires the right interval/retries/ceiling) by unwrapping the delegate. Two small parts:
(a) add a test-only accessor to ConnectionAwareRetryPolicy (plus the import org.apache.storm.shade.com.google.common.annotations.VisibleForTesting;):
@VisibleForTesting
RetryPolicy getDelegate() {
return delegate;
}(b) unwrap in CuratorUtilsTest.newCuratorUsesExponentialBackoffTest (add import org.apache.storm.shade.org.apache.curator.RetryPolicy; and assertTrue):
RetryPolicy retryPolicy = curator.getZookeeperClient().getRetryPolicy();
assertTrue(retryPolicy instanceof ConnectionAwareRetryPolicy,
"newCurator should wrap the retry policy in a ConnectionAwareRetryPolicy");
StormBoundedExponentialBackoffRetry policy =
(StormBoundedExponentialBackoffRetry) ((ConnectionAwareRetryPolicy) retryPolicy).getDelegate();
// existing getBaseSleepTimeMs / getN / getSleepTimeMs assertions stay as-is2. No unit tests for the new logic
ConnectionAwareRetryPolicy is 125 lines of connection-state/retry logic with no coverage (the PR's validation is a production anecdote). I wrote two classes — both pass locally (7/7 and 2/2). The policy tests use Mockito by capturing the ConnectionStateListener from bind(), so no real ZooKeeper is needed.
ConnectionAwareRetryPolicyTest.java (storm-client, 7 cases)
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.utils;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import org.apache.storm.shade.org.apache.curator.RetryPolicy;
import org.apache.storm.shade.org.apache.curator.RetrySleeper;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
import org.apache.storm.shade.org.apache.curator.framework.listen.Listenable;
import org.apache.storm.shade.org.apache.curator.framework.state.ConnectionState;
import org.apache.storm.shade.org.apache.curator.framework.state.ConnectionStateListener;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
/**
* Unit tests for {@link ConnectionAwareRetryPolicy}. The ZooKeeper connection state is driven by
* capturing the {@link ConnectionStateListener} that {@link ConnectionAwareRetryPolicy#bind} registers
* and firing state transitions at it, so no real ZooKeeper is required.
*/
public class ConnectionAwareRetryPolicyTest {
private static final int SESSION_TIMEOUT_MS = 20_000;
// JUnit 5 creates a fresh test instance per method, so these are clean for each test.
private final RetryPolicy delegate = mock(RetryPolicy.class);
private final RetrySleeper sleeper = mock(RetrySleeper.class);
private final CuratorFramework zk = mock(CuratorFramework.class);
private ConnectionStateListener listener;
/**
* Builds the policy, binds it to the mock framework, and captures the registered listener.
*
* @param zkSupplier what {@code zkSupplier.get()} returns inside allowRetry (the real wiring passes
* {@code zkRef::get}; pass {@code () -> null} to exercise the not-yet-bound path)
*/
@SuppressWarnings("unchecked")
private ConnectionAwareRetryPolicy build(Supplier<CuratorFramework> zkSupplier) {
Listenable<ConnectionStateListener> listenable = mock(Listenable.class);
when(zk.getConnectionStateListenable()).thenReturn(listenable);
ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS);
policy.bind(zk);
ArgumentCaptor<ConnectionStateListener> captor = ArgumentCaptor.forClass(ConnectionStateListener.class);
verify(listenable).addListener(captor.capture());
listener = captor.getValue();
return policy;
}
private void fire(ConnectionState state) {
listener.stateChanged(zk, state);
}
@Test
public void connectedStateDelegatesToBackoffPolicy() {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
// Default state is CONNECTED (no event fired).
boolean result = policy.allowRetry(0, 0L, sleeper);
assertTrue(result, "when connected, must propagate the delegate's decision");
verify(delegate).allowRetry(0, 0L, sleeper);
}
@Test
public void reconnectedStateDelegatesToBackoffPolicy() {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
fire(ConnectionState.RECONNECTED);
boolean result = policy.allowRetry(2, 100L, sleeper);
assertFalse(result, "RECONNECTED is a healthy state and must delegate");
verify(delegate).allowRetry(2, 100L, sleeper);
}
@Test
public void suspendedBlocksUntilConnectedThenRetriesImmediately() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
fire(ConnectionState.SUSPENDED);
boolean result = policy.allowRetry(0, 0L, sleeper);
assertTrue(result, "should retry once the SendThread has reconnected");
verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
verifyNoInteractions(sleeper); // retries immediately, no backoff sleep
verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
public void lostBlocksUntilConnectedThenRetriesImmediately() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
fire(ConnectionState.LOST);
boolean result = policy.allowRetry(3, 1234L, sleeper);
assertTrue(result, "LOST should also wait for reconnection then retry");
verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(false);
fire(ConnectionState.SUSPENDED);
boolean result = policy.allowRetry(0, 0L, sleeper);
assertFalse(result, "should abandon when not reconnected within the session timeout");
verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
public void interruptedWhileWaitingReturnsFalseAndPreservesInterruptFlag() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(zk.blockUntilConnected(anyInt(), any())).thenThrow(new InterruptedException("test"));
fire(ConnectionState.SUSPENDED);
boolean result = policy.allowRetry(0, 0L, sleeper);
assertFalse(result, "an interrupt while waiting should abandon the retry");
assertTrue(Thread.interrupted(), "interrupt flag must be preserved (this check also clears it)");
}
@Test
public void suspendedWithNullFrameworkFallsThroughToDelegate() {
// zkSupplier returns null (framework not yet available), even though bind() ran on the mock.
ConnectionAwareRetryPolicy policy = build(() -> null);
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
fire(ConnectionState.SUSPENDED);
boolean result = policy.allowRetry(1, 50L, sleeper);
assertTrue(result, "with no framework available yet, must fall through to the delegate");
verify(delegate).allowRetry(1, 50L, sleeper);
}
}SupervisorHeartbeatTest.java (storm-server, 2 cases)
/*
* Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
* 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
* and limitations under the License.
*/
package org.apache.storm.daemon.supervisor.timer;
import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.storm.cluster.IStormClusterState;
import org.apache.storm.daemon.supervisor.Supervisor;
import org.apache.storm.generated.SupervisorInfo;
import org.apache.storm.scheduler.ISupervisor;
import org.apache.storm.utils.Utils;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
public class SupervisorHeartbeatTest {
private static final String SUPERVISOR_ID = "supervisor-1";
/**
* Mocks a Supervisor with just enough state for {@code buildSupervisorInfo} to produce one
* SupervisorInfo entry (the non-NUMA branch), so {@code run()} reaches the ZK heartbeat write.
*/
private Supervisor mockSupervisor(IStormClusterState clusterState) {
Supervisor supervisor = mock(Supervisor.class);
ISupervisor iSupervisor = mock(ISupervisor.class);
when(iSupervisor.getMetadata()).thenReturn(Arrays.asList(6700, 6701));
when(supervisor.getStormClusterState()).thenReturn(clusterState);
when(supervisor.getId()).thenReturn(SUPERVISOR_ID);
when(supervisor.getiSupervisor()).thenReturn(iSupervisor);
when(supervisor.getCurrAssignment()).thenReturn(new AtomicReference<>(new HashMap<>()));
when(supervisor.getHostName()).thenReturn("host-1");
when(supervisor.getAssignmentId()).thenReturn(SUPERVISOR_ID);
when(supervisor.getThriftServerPort()).thenReturn(6627);
when(supervisor.getUpTime()).thenReturn(Utils.makeUptimeComputer());
when(supervisor.getStormVersion()).thenReturn("test-version");
return supervisor;
}
private Map<String, Object> baseConf() {
return new HashMap<>(Utils.readStormConfig());
}
@Test
public void transientHeartbeatFailureIsSwallowedInsteadOfKillingTheProcess() {
IStormClusterState clusterState = mock(IStormClusterState.class);
// Simulate a transient ZK write failure during the heartbeat cycle (the real one is a wrapped
// KeeperException$ConnectionLoss; the catch is on Exception, so a RuntimeException exercises the same path).
doThrow(new RuntimeException("simulated transient ZK failure"))
.when(clusterState).supervisorHeartbeat(anyString(), any(SupervisorInfo.class));
Supervisor supervisor = mockSupervisor(clusterState);
SupervisorHeartbeat heartbeat = new SupervisorHeartbeat(baseConf(), supervisor);
// The whole point of the fix: run() must not propagate — otherwise the StormTimer's
// DefaultUncaughtExceptionHandler would exit the supervisor process.
assertDoesNotThrow(heartbeat::run);
// It must have actually attempted the write (i.e. reached the failing ZK call).
verify(clusterState, atLeastOnce()).supervisorHeartbeat(anyString(), any(SupervisorInfo.class));
}
@Test
public void successfulCycleSendsTheHeartbeat() {
IStormClusterState clusterState = mock(IStormClusterState.class);
Supervisor supervisor = mockSupervisor(clusterState);
SupervisorHeartbeat heartbeat = new SupervisorHeartbeat(baseConf(), supervisor);
assertDoesNotThrow(heartbeat::run);
verify(clusterState).supervisorHeartbeat(eq(SUPERVISOR_ID), any(SupervisorInfo.class));
}
}3. Blast radius — please confirm intent
newCurator() backs every Curator client (Nimbus, blobstore, Trident TransactionalState, ClientZookeeper), so this changes ZK retry behavior cluster-wide, not just the supervisor heartbeat. On SUSPENDED/LOST, operations that previously failed fast can now block up to storm.zookeeper.session.timeout (20s default). Is that acceptable for the Nimbus/leader paths, which weren't exercised here?
4. No overall bound on the SUSPENDED/LOST path (optional)
That branch ignores retryCount/elapsedTimeMs and the delegate's retry.times cap, so a flapping connection can retry indefinitely and wedge the calling thread (e.g. the heartbeat timer). Consider still consulting the delegate's budget before waiting, e.g.:
// honour the configured retry budget even while suspended
if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
return false;
}5. Thundering herd on reconnect (optional)
The deliberate immediate retry (no jitter) means many suspended clients retry in lockstep the instant the ensemble member recovers; the old backoff spread them out. A little jitter before return true; would help on large clusters.
6 & 7. The heartbeat catch is broader than the bug (see inline suggestion)
It wraps getNumaMap()/buildSupervisorInfo() (local, non-ZK logic) too, so a persistent local misconfig (e.g. a malformed supervisor.numa.meta) would be swallowed every cycle forever instead of surfacing. The inline suggestion scopes the catch to just the ZK write — and also corrects the comment: nimbus.supervisor.timeout.secs defaults to 60s, not 30s (~12 tolerated misses at a 5s frequency, not 6).
| try { | ||
| Map<String, Object> validatedNumaMap = SupervisorUtils.getNumaMap(conf); | ||
| Map<String, SupervisorInfo> supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap); | ||
| for (Map.Entry<String, SupervisorInfo> supervisorInfoEntry : supervisorInfoList.entrySet()) { | ||
| stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue()); | ||
| } | ||
| } catch (Exception e) { | ||
| // A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 30s) | ||
| // allows 6 missed beats before the supervisor is considered dead. The next | ||
| // heartbeat (5s later) will retry. Killing the entire supervisor process | ||
| // over one missed beat is not acceptable. | ||
| LOG.warn("Supervisor heartbeat failed, will retry next cycle", e); | ||
| } |
There was a problem hiding this comment.
Scope the catch to the ZK write so a transient blip is skipped, but a persistent local failure (NUMA map, resource capacities) still surfaces instead of being swallowed forever. Also fixes the default in the comment (60s, not 30s).
| try { | |
| Map<String, Object> validatedNumaMap = SupervisorUtils.getNumaMap(conf); | |
| Map<String, SupervisorInfo> supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap); | |
| for (Map.Entry<String, SupervisorInfo> supervisorInfoEntry : supervisorInfoList.entrySet()) { | |
| stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue()); | |
| } | |
| } catch (Exception e) { | |
| // A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 30s) | |
| // allows 6 missed beats before the supervisor is considered dead. The next | |
| // heartbeat (5s later) will retry. Killing the entire supervisor process | |
| // over one missed beat is not acceptable. | |
| LOG.warn("Supervisor heartbeat failed, will retry next cycle", e); | |
| } | |
| Map<String, Object> validatedNumaMap = SupervisorUtils.getNumaMap(conf); | |
| Map<String, SupervisorInfo> supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap); | |
| for (Map.Entry<String, SupervisorInfo> supervisorInfoEntry : supervisorInfoList.entrySet()) { | |
| try { | |
| stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue()); | |
| } catch (Exception e) { | |
| // A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 60s) tolerates | |
| // several missed beats; the next beat (5s later) retries. Local/config errors are left | |
| // outside this catch so a persistent misconfiguration still surfaces instead of being | |
| // swallowed every cycle. | |
| LOG.warn("Supervisor heartbeat for {} failed, will retry next cycle", | |
| supervisorInfoEntry.getKey(), e); | |
| } | |
| } |
|
@rzo1 @GGraziadei the point highlighted above is more of a design decision. I'm inclined to say that yes, we want all processes to be resilient to zookeeper connection blips. What is your take on this?
|
|
Please check the build failures. Thanks :) |
…catch - Fix CuratorUtilsTest ClassCastException (unwrap ConnectionAwareRetryPolicy) - Add @VisibleForTesting getDelegate() accessor - Add ConnectionAwareRetryPolicyTest (8 cases) and SupervisorHeartbeatTest (2 cases) - Add delegate budget check before blocking on SUSPENDED/LOST - Add RECONNECT_JITTER_MS constant (100ms) to avoid thundering herd - Scope SupervisorHeartbeat catch to ZK write only (not getNumaMap/buildSupervisorInfo) - Fix heartbeat comment (liveness via ephemeral ZK nodes, not heartbeat timestamps)
reiabreu
left a comment
There was a problem hiding this comment.
This review was generated with the help of an LLM (Claude) and reviewed by me before posting.
Follow-up on the retry-bound change (my earlier point 4). The current form opens the SUSPENDED/LOST branch with if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) return false; to honour the budget — but delegate.allowRetry (Curator's SleepingRetry) sleeps the exponential backoff as a side effect before returning true. So the suspended path now does backoff sleep → blockUntilConnected() → jitter, which reintroduces exactly the blind backoff this policy was built to avoid (the old "retry immediately" comment was dropped). The unit tests don't catch it because delegate is mocked, so the real sleep never runs.
Proposed fix: bound the suspended path by retry count (a pure check, no sleep) instead of calling the sleeping delegate. maxRetries is wired from the same storm.zookeeper.retry.times the delegate already uses, so both paths give up after the same number of attempts; the bounded jitter stays. The healthy (CONNECTED/RECONNECTED) path is unchanged and still uses the real backoff, which is correct there.
One API note: this adds a 4-arg constructor to ConnectionAwareRetryPolicy (your class) — (..., int sessionTimeoutMs, int maxRetries).
I ran it locally: ConnectionAwareRetryPolicyTest 8/8 and CuratorUtilsTest 5/5 pass (the only build failure is the unrelated python3-test CLI harness, which fails the same way without this change). The tests now assert verify(delegate, never()).allowRetry(...) on the suspended path — i.e. no blind backoff — plus a bounded-jitter check, and drive the budget by retry count.
The core production hunk is inline below as a suggestion; it depends on the constructor + CuratorUtils edits, so apply it together with the full patch:
diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
index eb041bdef..44d5a20fa 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
@@ -60,6 +60,7 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
private final RetryPolicy delegate;
private final Supplier<CuratorFramework> zkSupplier;
private final int sessionTimeoutMs;
+ private final int maxRetries;
private final AtomicReference<ConnectionState> connectionState =
new AtomicReference<>(ConnectionState.CONNECTED);
@@ -69,13 +70,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
* the framework does not exist until {@code builder.build()} returns
* @param sessionTimeoutMs upper bound (in ms) for waiting on reconnection; typically
* {@code storm.zookeeper.session.timeout}
+ * @param maxRetries cap on the number of retries on the SUSPENDED/LOST path, applied without a
+ * backoff sleep. Should match the delegate's retry budget
+ * ({@code storm.zookeeper.retry.times}) so the two paths give up after the
+ * same number of attempts.
*/
public ConnectionAwareRetryPolicy(RetryPolicy delegate,
Supplier<CuratorFramework> zkSupplier,
- int sessionTimeoutMs) {
+ int sessionTimeoutMs,
+ int maxRetries) {
this.delegate = delegate;
this.zkSupplier = zkSupplier;
this.sessionTimeoutMs = sessionTimeoutMs;
+ this.maxRetries = maxRetries;
}
/**
@@ -104,14 +111,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
ConnectionState state = connectionState.get();
if (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST) {
- // Honour the configured retry budget even while suspended.
- if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
+ // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
+ // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
+ // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
+ // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
+ if (retryCount >= maxRetries) {
+ LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
+ state, maxRetries, retryCount);
return false;
}
CuratorFramework zk = zkSupplier.get();
if (zk == null) {
- // Framework not yet available; delegate already approved the retry
+ // Framework not yet available; within budget, so allow the retry.
return true;
}
diff --git a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
index d2365bd1a..d15e2b143 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
@@ -63,13 +63,15 @@ public class CuratorUtils {
// Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST.
AtomicReference<CuratorFramework> zkRef = new AtomicReference<>();
int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT));
+ int retryTimes = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES));
ConnectionAwareRetryPolicy connectionAwarePolicy = new ConnectionAwareRetryPolicy(
new StormBoundedExponentialBackoffRetry(
ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL)),
ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL_CEILING)),
- ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES))),
+ retryTimes),
zkRef::get,
- sessionTimeoutMs);
+ sessionTimeoutMs,
+ retryTimes);
builder.retryPolicy(connectionAwarePolicy);
if (defaultAcl != null) {
Full patch (production + updated tests)
diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
index eb041bdef..44d5a20fa 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
@@ -60,6 +60,7 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
private final RetryPolicy delegate;
private final Supplier<CuratorFramework> zkSupplier;
private final int sessionTimeoutMs;
+ private final int maxRetries;
private final AtomicReference<ConnectionState> connectionState =
new AtomicReference<>(ConnectionState.CONNECTED);
@@ -69,13 +70,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
* the framework does not exist until {@code builder.build()} returns
* @param sessionTimeoutMs upper bound (in ms) for waiting on reconnection; typically
* {@code storm.zookeeper.session.timeout}
+ * @param maxRetries cap on the number of retries on the SUSPENDED/LOST path, applied without a
+ * backoff sleep. Should match the delegate's retry budget
+ * ({@code storm.zookeeper.retry.times}) so the two paths give up after the
+ * same number of attempts.
*/
public ConnectionAwareRetryPolicy(RetryPolicy delegate,
Supplier<CuratorFramework> zkSupplier,
- int sessionTimeoutMs) {
+ int sessionTimeoutMs,
+ int maxRetries) {
this.delegate = delegate;
this.zkSupplier = zkSupplier;
this.sessionTimeoutMs = sessionTimeoutMs;
+ this.maxRetries = maxRetries;
}
/**
@@ -104,14 +111,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
ConnectionState state = connectionState.get();
if (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST) {
- // Honour the configured retry budget even while suspended.
- if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
+ // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
+ // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
+ // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
+ // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
+ if (retryCount >= maxRetries) {
+ LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
+ state, maxRetries, retryCount);
return false;
}
CuratorFramework zk = zkSupplier.get();
if (zk == null) {
- // Framework not yet available; delegate already approved the retry
+ // Framework not yet available; within budget, so allow the retry.
return true;
}
diff --git a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
index d2365bd1a..d15e2b143 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
@@ -63,13 +63,15 @@ public class CuratorUtils {
// Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST.
AtomicReference<CuratorFramework> zkRef = new AtomicReference<>();
int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT));
+ int retryTimes = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES));
ConnectionAwareRetryPolicy connectionAwarePolicy = new ConnectionAwareRetryPolicy(
new StormBoundedExponentialBackoffRetry(
ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL)),
ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL_CEILING)),
- ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES))),
+ retryTimes),
zkRef::get,
- sessionTimeoutMs);
+ sessionTimeoutMs,
+ retryTimes);
builder.retryPolicy(connectionAwarePolicy);
if (defaultAcl != null) {
diff --git a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
index 38236fffa..b6259610c 100644
--- a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
+++ b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
@@ -34,6 +34,9 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.longThat;
+import static org.mockito.Mockito.atMost;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
@@ -47,6 +50,9 @@ import static org.mockito.Mockito.when;
public class ConnectionAwareRetryPolicyTest {
private static final int SESSION_TIMEOUT_MS = 20_000;
+ private static final int MAX_RETRIES = 5;
+ // Must stay in sync with ConnectionAwareRetryPolicy's RECONNECT_JITTER_MS bound.
+ private static final long JITTER_BOUND_MS = 100L;
private final RetryPolicy delegate = mock(RetryPolicy.class);
private final RetrySleeper sleeper = mock(RetrySleeper.class);
@@ -58,7 +64,8 @@ public class ConnectionAwareRetryPolicyTest {
Listenable<ConnectionStateListener> listenable = mock(Listenable.class);
when(zk.getConnectionStateListenable()).thenReturn(listenable);
- ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS);
+ ConnectionAwareRetryPolicy policy =
+ new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS, MAX_RETRIES);
policy.bind(zk);
ArgumentCaptor<ConnectionStateListener> captor = ArgumentCaptor.forClass(ConnectionStateListener.class);
@@ -76,7 +83,8 @@ public class ConnectionAwareRetryPolicyTest {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
- // Default state is CONNECTED (no event fired).
+ // Default state is CONNECTED (no event fired). The backoff policy is correct here: a
+ // retryable error while connected is not a disconnect, so normal exponential backoff applies.
boolean result = policy.allowRetry(0, 0L, sleeper);
assertTrue(result, "when connected, must propagate the delegate's decision");
@@ -96,60 +104,64 @@ public class ConnectionAwareRetryPolicyTest {
}
@Test
- public void suspendedBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+ public void suspendedWaitsForReconnectWithoutBlindBackoffThenRetries() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(0, 0L, sleeper);
+ boolean result = policy.allowRetry(0, 0L, sleeper); // retryCount 0 < MAX_RETRIES
assertTrue(result, "should retry once the SendThread has reconnected");
verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+ // The whole point of this policy: on SUSPENDED/LOST it must NOT invoke the delegate's
+ // exponential backoff sleep. Reaching the delegate here would reintroduce the blind backoff.
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
+ // Only bounded reconnection jitter may sleep — never a multi-second backoff.
+ verify(sleeper, atMost(1)).sleepFor(longThat(ms -> ms >= 0 && ms < JITTER_BOUND_MS),
+ eq(TimeUnit.MILLISECONDS));
}
@Test
- public void lostBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+ public void lostWaitsForReconnectWithoutBlindBackoffThenRetries() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
fire(ConnectionState.LOST);
- boolean result = policy.allowRetry(3, 1234L, sleeper);
+ boolean result = policy.allowRetry(3, 1234L, sleeper); // 3 < MAX_RETRIES
assertTrue(result, "LOST should also wait for reconnection then retry");
verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(false);
fire(ConnectionState.SUSPENDED);
boolean result = policy.allowRetry(0, 0L, sleeper);
assertFalse(result, "should abandon when not reconnected within the session timeout");
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
- public void suspendedHonoursDelegateRetryBudget() throws Exception {
+ public void suspendedAbandonsWhenRetryBudgetExhausted() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- // Delegate says no more retries allowed (budget exhausted)
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(100, 99999L, sleeper);
+ // retryCount == MAX_RETRIES → budget exhausted; must give up without waiting or backing off.
+ boolean result = policy.allowRetry(MAX_RETRIES, 99999L, sleeper);
- assertFalse(result, "should abandon when the delegate's retry budget is exhausted");
+ assertFalse(result, "should abandon when the retry budget is exhausted");
verify(zk, never()).blockUntilConnected(anyInt(), any());
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
public void interruptedWhileWaitingReturnsFalseAndPreservesInterruptFlag() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(anyInt(), any())).thenThrow(new InterruptedException("test"));
fire(ConnectionState.SUSPENDED);
@@ -160,15 +172,15 @@ public class ConnectionAwareRetryPolicyTest {
}
@Test
- public void suspendedWithNullFrameworkFallsThroughToDelegate() {
- // zkSupplier returns null (framework not yet available), even though bind() ran on the mock.
+ public void suspendedWithinBudgetButNullFrameworkAllowsRetry() {
+ // zkSupplier returns null (framework not yet bound), even though bind() ran on the mock.
ConnectionAwareRetryPolicy policy = build(() -> null);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(1, 50L, sleeper);
+ boolean result = policy.allowRetry(1, 50L, sleeper); // 1 < MAX_RETRIES
- assertTrue(result, "with no framework available yet, must fall through to the delegate");
- verify(delegate).allowRetry(1, 50L, sleeper);
+ assertTrue(result, "within budget but framework not available yet → allow the retry");
+ // Still no blind backoff on the suspended path.
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
}
(Separately from this — the blast-radius design question still stands for @rzo1 / @GGraziadei.)
| // Honour the configured retry budget even while suspended. | ||
| if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) { | ||
| return false; | ||
| } | ||
|
|
||
| CuratorFramework zk = zkSupplier.get(); | ||
| if (zk == null) { | ||
| // Framework not yet available; delegate already approved the retry | ||
| return true; | ||
| } |
There was a problem hiding this comment.
Core change: replace the budget check that calls the sleeping delegate with a pure retry-count cap, so the suspended path no longer reintroduces the exponential backoff sleep. Depends on the new maxRetries field/constructor arg and the CuratorUtils wiring in the patch above — apply those together or this won't compile.
| // Honour the configured retry budget even while suspended. | |
| if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) { | |
| return false; | |
| } | |
| CuratorFramework zk = zkSupplier.get(); | |
| if (zk == null) { | |
| // Framework not yet available; delegate already approved the retry | |
| return true; | |
| } | |
| // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On | |
| // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected() | |
| // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff | |
| // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid. | |
| if (retryCount >= maxRetries) { | |
| LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry", | |
| state, maxRetries, retryCount); | |
| return false; | |
| } | |
| CuratorFramework zk = zkSupplier.get(); | |
| if (zk == null) { | |
| // Framework not yet available; within budget, so allow the retry. | |
| return true; | |
| } |
|
Thanks @jkrauss82 for the report and fix, and @reiabreu for the review. I would leave the retry behaviour of the other Curator clients unchanged in this PR.
I'd keep this PR focused on the @jkrauss82, could you share the full stack trace for the Happy to reconsider if the stack trace shows something different. What do you think? |
|
thanks for looking into this @GGraziadei I don't have enough knowledge of the source and the actual inner working of the various components using Curator to contribute meaningful to the discussion regarding the scope of the PR. but I hope the logs below provide further insight. happy to provide more if necessary. also sharing this summary my local AI gave me regarding the log content, just to provide further context as it includes the various relevant timestamps:
|
|
Written with Claude (an LLM). Claude ran all the tests described below and drafted this comment; I reviewed it before posting. Summary: Claude's failover tests found that Problem. For background (
A forced version is deterministic: Fix for this PR. Wait for the reconnect only on foreground retries and delegate background retries to the wrapped policy. Curator passes private static boolean isForegroundRetry(RetrySleeper sleeper) {
return sleeper == RetryLoop.getDefaultRetrySleeper();
}
// in allowRetry(...):
if (isForegroundRetry(sleepSleeper) && (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST)) {
... // wait for reconnect, as before
}
return delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper);The patch below also bounds the suspended path by retry count instead of calling the delegate (which sleeps its backoff first), wires Trade-off. Background operations no longer survive a long outage. In an 8 s outage about 3,000 of 5,000 background reads gave up with the fix, the same as stock. The policy as written lost none, but only by blocking the event thread. Limits.
Re @GGraziadei's question (was the policy consulted?). The posted trace can't tell. Patch (on top of ab67921)diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
index eb041bdef..e3f90eb6c 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
@@ -23,6 +23,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Supplier;
import org.apache.storm.shade.com.google.common.annotations.VisibleForTesting;
+import org.apache.storm.shade.org.apache.curator.RetryLoop;
import org.apache.storm.shade.org.apache.curator.RetryPolicy;
import org.apache.storm.shade.org.apache.curator.RetrySleeper;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
@@ -44,7 +45,13 @@ import org.slf4j.LoggerFactory;
* the operation on the new connection. If the connection cannot be
* re-established within the session timeout, the retry is abandoned.
*
- * <p>For all other connection states, this policy delegates to the wrapped
+ * <p>Waiting is only done for <em>foreground</em> retries, which run on the caller's own thread.
+ * Curator also retries <em>background</em> ({@code inBackground}) operations, and there the policy is
+ * invoked from ZooKeeper's single event thread or from Curator's background thread. Blocking the event
+ * thread delays delivery of the very "connected" event that {@code blockUntilConnected} is waiting for, so
+ * background retries always delegate to the wrapped policy, whose sleeper only schedules a re-queue.
+ *
+ * <p>For all other connection states, and for background retries, this policy delegates to the wrapped
* delegate policy (typically a {@link StormBoundedExponentialBackoffRetry}).
*/
public class ConnectionAwareRetryPolicy implements RetryPolicy {
@@ -60,6 +67,7 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
private final RetryPolicy delegate;
private final Supplier<CuratorFramework> zkSupplier;
private final int sessionTimeoutMs;
+ private final int maxRetries;
private final AtomicReference<ConnectionState> connectionState =
new AtomicReference<>(ConnectionState.CONNECTED);
@@ -69,13 +77,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
* the framework does not exist until {@code builder.build()} returns
* @param sessionTimeoutMs upper bound (in ms) for waiting on reconnection; typically
* {@code storm.zookeeper.session.timeout}
+ * @param maxRetries cap on the number of retries on the SUSPENDED/LOST path, applied without a
+ * backoff sleep. Should match the delegate's retry budget
+ * ({@code storm.zookeeper.retry.times}) so the two paths give up after the
+ * same number of attempts.
*/
public ConnectionAwareRetryPolicy(RetryPolicy delegate,
Supplier<CuratorFramework> zkSupplier,
- int sessionTimeoutMs) {
+ int sessionTimeoutMs,
+ int maxRetries) {
this.delegate = delegate;
this.zkSupplier = zkSupplier;
this.sessionTimeoutMs = sessionTimeoutMs;
+ this.maxRetries = maxRetries;
}
/**
@@ -99,19 +113,33 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
return delegate;
}
+ /**
+ * Curator's foreground retry loop passes {@link RetryLoop#getDefaultRetrySleeper()} (a real sleep on the
+ * caller's thread). Background operations pass the operation itself, whose {@code sleepFor} only records a
+ * re-queue time, and whose retry runs on a thread that must not block.
+ */
+ private static boolean isForegroundRetry(RetrySleeper sleeper) {
+ return sleeper == RetryLoop.getDefaultRetrySleeper();
+ }
+
@Override
public boolean allowRetry(int retryCount, long elapsedTimeMs, RetrySleeper sleepSleeper) {
ConnectionState state = connectionState.get();
- if (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST) {
- // Honour the configured retry budget even while suspended.
- if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
+ if (isForegroundRetry(sleepSleeper) && (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST)) {
+ // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
+ // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
+ // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
+ // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
+ if (retryCount >= maxRetries) {
+ LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
+ state, maxRetries, retryCount);
return false;
}
CuratorFramework zk = zkSupplier.get();
if (zk == null) {
- // Framework not yet available; delegate already approved the retry
+ // Framework not yet available; within budget, so allow the retry.
return true;
}
diff --git a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
index d2365bd1a..d15e2b143 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
@@ -63,13 +63,15 @@ public class CuratorUtils {
// Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST.
AtomicReference<CuratorFramework> zkRef = new AtomicReference<>();
int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT));
+ int retryTimes = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES));
ConnectionAwareRetryPolicy connectionAwarePolicy = new ConnectionAwareRetryPolicy(
new StormBoundedExponentialBackoffRetry(
ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL)),
ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL_CEILING)),
- ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES))),
+ retryTimes),
zkRef::get,
- sessionTimeoutMs);
+ sessionTimeoutMs,
+ retryTimes);
builder.retryPolicy(connectionAwarePolicy);
if (defaultAcl != null) {
diff --git a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
index 38236fffa..5f50cfc66 100644
--- a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
+++ b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
@@ -20,6 +20,7 @@ package org.apache.storm.utils;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
+import org.apache.storm.shade.org.apache.curator.RetryLoop;
import org.apache.storm.shade.org.apache.curator.RetryPolicy;
import org.apache.storm.shade.org.apache.curator.RetrySleeper;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
@@ -43,13 +44,19 @@ import static org.mockito.Mockito.when;
* Unit tests for {@link ConnectionAwareRetryPolicy}. The ZooKeeper connection state is driven by
* capturing the {@link ConnectionStateListener} that {@link ConnectionAwareRetryPolicy#bind} registers
* and firing state transitions at it, so no real ZooKeeper is required.
+ *
+ * <p>Curator passes {@link RetryLoop#getDefaultRetrySleeper()} for foreground retries and the operation
+ * itself for background retries. The policy waits only for the former, so the tests use the real default
+ * sleeper to exercise the waiting path and a mock sleeper to stand in for a background retry.
*/
public class ConnectionAwareRetryPolicyTest {
private static final int SESSION_TIMEOUT_MS = 20_000;
+ private static final int MAX_RETRIES = 5;
private final RetryPolicy delegate = mock(RetryPolicy.class);
- private final RetrySleeper sleeper = mock(RetrySleeper.class);
+ private final RetrySleeper foregroundSleeper = RetryLoop.getDefaultRetrySleeper();
+ private final RetrySleeper backgroundSleeper = mock(RetrySleeper.class);
private final CuratorFramework zk = mock(CuratorFramework.class);
private ConnectionStateListener listener;
@@ -58,7 +65,8 @@ public class ConnectionAwareRetryPolicyTest {
Listenable<ConnectionStateListener> listenable = mock(Listenable.class);
when(zk.getConnectionStateListenable()).thenReturn(listenable);
- ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS);
+ ConnectionAwareRetryPolicy policy =
+ new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS, MAX_RETRIES);
policy.bind(zk);
ArgumentCaptor<ConnectionStateListener> captor = ArgumentCaptor.forClass(ConnectionStateListener.class);
@@ -76,11 +84,12 @@ public class ConnectionAwareRetryPolicyTest {
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
- // Default state is CONNECTED (no event fired).
- boolean result = policy.allowRetry(0, 0L, sleeper);
+ // Default state is CONNECTED (no event fired). The backoff policy is correct here: a
+ // retryable error while connected is not a disconnect, so normal exponential backoff applies.
+ boolean result = policy.allowRetry(0, 0L, foregroundSleeper);
assertTrue(result, "when connected, must propagate the delegate's decision");
- verify(delegate).allowRetry(0, 0L, sleeper);
+ verify(delegate).allowRetry(0, 0L, foregroundSleeper);
}
@Test
@@ -89,86 +98,116 @@ public class ConnectionAwareRetryPolicyTest {
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
fire(ConnectionState.RECONNECTED);
- boolean result = policy.allowRetry(2, 100L, sleeper);
+ boolean result = policy.allowRetry(2, 100L, foregroundSleeper);
assertFalse(result, "RECONNECTED is a healthy state and must delegate");
- verify(delegate).allowRetry(2, 100L, sleeper);
+ verify(delegate).allowRetry(2, 100L, foregroundSleeper);
}
@Test
- public void suspendedBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+ public void suspendedForegroundRetryWaitsForReconnectWithoutBlindBackoff() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(0, 0L, sleeper);
+ boolean result = policy.allowRetry(0, 0L, foregroundSleeper); // retryCount 0 < MAX_RETRIES
assertTrue(result, "should retry once the SendThread has reconnected");
verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+ // On SUSPENDED/LOST the delegate's exponential backoff sleep must not run: reaching it here
+ // would reintroduce the blind backoff this policy exists to avoid.
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
- public void lostBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+ public void lostForegroundRetryWaitsForReconnectWithoutBlindBackoff() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
fire(ConnectionState.LOST);
- boolean result = policy.allowRetry(3, 1234L, sleeper);
+ boolean result = policy.allowRetry(3, 1234L, foregroundSleeper); // 3 < MAX_RETRIES
assertTrue(result, "LOST should also wait for reconnection then retry");
verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
- public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
+ public void suspendedBackgroundRetryNeverBlocksAndDelegates() throws Exception {
+ // Background retries run on ZooKeeper's event thread (or Curator's background thread). Blocking there
+ // delays the "connected" event the wait is for, so they must go to the delegate, whose sleeper only
+ // schedules a re-queue.
ConnectionAwareRetryPolicy policy = build(() -> zk);
when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
+
+ fire(ConnectionState.SUSPENDED);
+ boolean result = policy.allowRetry(0, 0L, backgroundSleeper);
+
+ assertTrue(result, "a background retry must get the delegate's answer");
+ verify(delegate).allowRetry(0, 0L, backgroundSleeper);
+ verify(zk, never()).blockUntilConnected(anyInt(), any());
+ }
+
+ @Test
+ public void lostBackgroundRetryNeverBlocksAndDelegates() throws Exception {
+ ConnectionAwareRetryPolicy policy = build(() -> zk);
+ when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
+
+ fire(ConnectionState.LOST);
+ boolean result = policy.allowRetry(2, 50L, backgroundSleeper);
+
+ assertFalse(result, "a background retry must get the delegate's answer, including a refusal");
+ verify(delegate).allowRetry(2, 50L, backgroundSleeper);
+ verify(zk, never()).blockUntilConnected(anyInt(), any());
+ }
+
+ @Test
+ public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
+ ConnectionAwareRetryPolicy policy = build(() -> zk);
when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(false);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(0, 0L, sleeper);
+ boolean result = policy.allowRetry(0, 0L, foregroundSleeper);
assertFalse(result, "should abandon when not reconnected within the session timeout");
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
- public void suspendedHonoursDelegateRetryBudget() throws Exception {
+ public void suspendedAbandonsWhenRetryBudgetExhausted() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- // Delegate says no more retries allowed (budget exhausted)
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(100, 99999L, sleeper);
+ // retryCount == MAX_RETRIES → budget exhausted; must give up without waiting or backing off.
+ boolean result = policy.allowRetry(MAX_RETRIES, 99999L, foregroundSleeper);
- assertFalse(result, "should abandon when the delegate's retry budget is exhausted");
+ assertFalse(result, "should abandon when the retry budget is exhausted");
verify(zk, never()).blockUntilConnected(anyInt(), any());
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
@Test
public void interruptedWhileWaitingReturnsFalseAndPreservesInterruptFlag() throws Exception {
ConnectionAwareRetryPolicy policy = build(() -> zk);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
when(zk.blockUntilConnected(anyInt(), any())).thenThrow(new InterruptedException("test"));
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(0, 0L, sleeper);
+ boolean result = policy.allowRetry(0, 0L, foregroundSleeper);
assertFalse(result, "an interrupt while waiting should abandon the retry");
assertTrue(Thread.interrupted(), "interrupt flag must be preserved (this check also clears it)");
}
@Test
- public void suspendedWithNullFrameworkFallsThroughToDelegate() {
- // zkSupplier returns null (framework not yet available), even though bind() ran on the mock.
+ public void suspendedWithinBudgetButNullFrameworkAllowsRetry() {
+ // zkSupplier returns null (framework not yet bound), even though bind() ran on the mock.
ConnectionAwareRetryPolicy policy = build(() -> null);
- when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
fire(ConnectionState.SUSPENDED);
- boolean result = policy.allowRetry(1, 50L, sleeper);
+ boolean result = policy.allowRetry(1, 50L, foregroundSleeper); // 1 < MAX_RETRIES
- assertTrue(result, "with no framework available yet, must fall through to the delegate");
- verify(delegate).allowRetry(1, 50L, sleeper);
+ assertTrue(result, "within budget but framework not available yet → allow the retry");
+ // Still no blind backoff on the suspended path.
+ verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
}
}Failover test program (throwaway JUnit test, not for merging)package org.apache.storm.utils;
import java.io.PrintWriter;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.curator.test.TestingServer;
import org.apache.storm.shade.org.apache.curator.RetryPolicy;
import org.apache.storm.shade.org.apache.curator.RetrySleeper;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.storm.shade.org.apache.curator.framework.api.BackgroundCallback;
import org.apache.storm.shade.org.apache.curator.framework.state.ConnectionState;
import org.apache.storm.shade.org.apache.zookeeper.KeeperException;
import org.junit.jupiter.api.Test;
/**
* THROWAWAY experiment (not for commit): does ConnectionAwareRetryPolicy.allowRetry() ever block on ZooKeeper's
* event thread, and does that delay delivery of the RECONNECTED notification?
*/
public class ConnectionAwareRetryPolicyReproTest {
private static final String OUT = "/tmp/claude-1000/-home-rui-workspace/a5338b11-f59f-4b5e-aa07-12c2598cd7bd/scratchpad/storm9153/repro_out.txt";
private static final int SESSION_MS = 10_000;
private static final int CONN_MS = 5_000;
private static final int MAX_RETRIES = 5;
private static final int BURST = 5000;
private final List<String> out = Collections.synchronizedList(new ArrayList<>());
private long t0;
private void log(String s) {
out.add(String.format("%7d ms | %s", System.currentTimeMillis() - t0, s));
}
/** Records every allowRetry call: thread, the policy's own view of the state, duration, result. */
static class RecordingPolicy extends ConnectionAwareRetryPolicy {
final List<String> calls = Collections.synchronizedList(new ArrayList<>());
final AtomicInteger total = new AtomicInteger();
final AtomicInteger onEventThread = new AtomicInteger();
final AtomicInteger onEventThreadWhileDown = new AtomicInteger();
final AtomicInteger blockedOnEventThread = new AtomicInteger();
final AtomicInteger maxMs = new AtomicInteger();
final AtomicInteger offEventThread = new AtomicInteger();
final AtomicInteger offEventThreadWhileDown = new AtomicInteger();
final AtomicInteger maxMsOffEventWhileDown = new AtomicInteger();
final java.util.Set<String> offEventThreadNames = java.util.concurrent.ConcurrentHashMap.newKeySet();
final long start;
final Field stateField;
RecordingPolicy(RetryPolicy delegate, java.util.function.Supplier<CuratorFramework> sup, long start) throws Exception {
super(delegate, sup, SESSION_MS, MAX_RETRIES);
this.start = start;
stateField = ConnectionAwareRetryPolicy.class.getDeclaredField("connectionState");
stateField.setAccessible(true);
}
@Override
@SuppressWarnings("unchecked")
public boolean allowRetry(int retryCount, long elapsedTimeMs, RetrySleeper sleeper) {
String thread = Thread.currentThread().getName();
boolean event = thread.endsWith("-EventThread");
Object state;
try {
state = ((AtomicReference<ConnectionState>) stateField.get(this)).get();
} catch (Exception e) {
state = "?";
}
long a = System.currentTimeMillis();
boolean r = super.allowRetry(retryCount, elapsedTimeMs, sleeper);
long d = System.currentTimeMillis() - a;
total.incrementAndGet();
boolean down = state == ConnectionState.SUSPENDED || state == ConnectionState.LOST;
if (!event) {
offEventThread.incrementAndGet();
offEventThreadNames.add(thread);
if (down) {
offEventThreadWhileDown.incrementAndGet();
maxMsOffEventWhileDown.accumulateAndGet((int) d, Math::max);
}
}
if (event) {
onEventThread.incrementAndGet();
if (down) {
onEventThreadWhileDown.incrementAndGet();
}
if (d > 500) {
blockedOnEventThread.incrementAndGet();
}
}
maxMs.accumulateAndGet((int) d, Math::max);
if (down || d > 200) {
calls.add(String.format("%7d ms | thread=%s event=%s state=%s retry=%d sleeperClass=%s took=%dms -> %s",
System.currentTimeMillis() - start, thread, event, state, retryCount,
sleeper.getClass().getSimpleName(), d, r));
}
return r;
}
}
private void experiment(String label, boolean aware, int blips, long downMs, long upMs) throws Exception {
t0 = System.currentTimeMillis();
out.add("");
out.add("===== " + label + " | aware=" + aware + " blips=" + blips + " down=" + downMs + "ms up=" + upMs + "ms =====");
try (TestingServer server = new TestingServer(true)) {
AtomicReference<CuratorFramework> ref = new AtomicReference<>();
RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
RecordingPolicy rec = aware ? new RecordingPolicy(stock, ref::get, t0) : null;
CuratorFramework zk = CuratorFrameworkFactory.builder()
.connectString(server.getConnectString())
.sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS)
.retryPolicy(aware ? rec : stock).build();
ref.set(zk);
if (aware) {
rec.bind(zk);
}
final long[] lastRestartDone = {0};
final long[] reconnectedAt = {0};
zk.getConnectionStateListenable().addListener((c, s) -> {
log("state -> " + s);
if (s == ConnectionState.RECONNECTED) {
reconnectedAt[0] = System.currentTimeMillis();
}
});
zk.start();
zk.blockUntilConnected(10, TimeUnit.SECONDS);
zk.create().forPath("/r");
out.add("negotiated session timeout = " + zk.getZookeeperClient().getZooKeeper().getSessionTimeout() + " ms");
AtomicInteger connLoss = new AtomicInteger();
AtomicInteger ok = new AtomicInteger();
AtomicInteger other = new AtomicInteger();
final AtomicInteger done = new AtomicInteger();
final long[] allDoneAt = {0};
final long[] probeSentAt = {0};
final long[] probeDoneAt = {0};
BackgroundCallback cb = (c, e) -> {
int rc = e.getResultCode();
if (rc == KeeperException.Code.CONNECTIONLOSS.intValue()) {
connLoss.incrementAndGet();
} else if (rc == KeeperException.Code.OK.intValue()) {
ok.incrementAndGet();
} else {
other.incrementAndGet();
}
if (done.incrementAndGet() == BURST) {
allDoneAt[0] = System.currentTimeMillis();
}
};
for (int b = 0; b < blips; b++) {
reconnectedAt[0] = 0;
for (int i = 0; i < BURST; i++) {
zk.getData().inBackground(cb).forPath("/r"); // many in flight when the socket drops
}
log("blip " + b + ": stopping server");
server.stop();
if (downMs > 2000) {
Thread.sleep(downMs - 1000);
probeSentAt[0] = System.currentTimeMillis();
zk.getData().inBackground((c, e) -> probeDoneAt[0] = System.currentTimeMillis()).forPath("/r");
Thread.sleep(1000);
} else {
Thread.sleep(downMs);
}
server.restart();
lastRestartDone[0] = System.currentTimeMillis();
log("blip " + b + ": server restarted");
Thread.sleep(upMs);
}
long waitUntil = System.currentTimeMillis() + 40_000;
while (reconnectedAt[0] == 0 && System.currentTimeMillis() < waitUntil) {
Thread.sleep(50);
}
long doneDeadline = System.currentTimeMillis() + 40_000;
while (done.get() < BURST && System.currentTimeMillis() < doneDeadline) {
Thread.sleep(50);
}
Thread.sleep(500);
out.add("callbacks: CONNECTIONLOSS=" + connLoss + " OK=" + ok + " other=" + other + " (CONNECTIONLOSS = ops whose retries were EXHAUSTED; Curator hides retried failures from this callback)");
out.add("burst callbacks completed: " + done + "/" + BURST
+ (allDoneAt[0] == 0 ? " (some NEVER completed)" : "; last one " + (allDoneAt[0] - lastRestartDone[0]) + " ms after restart (negative = before)"));
if (probeSentAt[0] != 0) {
out.add("probe op sent ~1s before restart: " + (probeDoneAt[0] == 0 ? "NOT COMPLETED"
: "completed " + (probeDoneAt[0] - lastRestartDone[0]) + " ms after restart"));
}
out.add("last restart -> RECONNECTED notification latency = "
+ (reconnectedAt[0] == 0 ? "NEVER (40s)" : (reconnectedAt[0] - lastRestartDone[0]) + " ms"));
if (aware) {
out.add("policy.allowRetry calls: total=" + rec.total + " onEventThread=" + rec.onEventThread
+ " onEventThreadWhileSuspendedOrLost=" + rec.onEventThreadWhileDown
+ " blockedOver500msOnEventThread=" + rec.blockedOnEventThread + " maxCallMs=" + rec.maxMs);
out.add(" OFF event thread: calls=" + rec.offEventThread + " whileSuspendedOrLost=" + rec.offEventThreadWhileDown
+ " longestWhileDownMs=" + rec.maxMsOffEventWhileDown + " threads=" + rec.offEventThreadNames);
synchronized (rec.calls) {
int n = 0;
for (String c : rec.calls) {
if (n++ < 12) {
out.add(" " + c);
}
}
if (rec.calls.size() > 12) {
out.add(" ... (" + rec.calls.size() + " notable calls total)");
}
}
}
zk.close();
}
}
/**
* Forces the suspected situation: the policy believes SUSPENDED, Curator is disconnected, and a failed-operation
* callback runs on ZooKeeper's event thread and calls allowRetry(). Does that stall reconnect delivery?
*/
private void mechanism() throws Exception {
t0 = System.currentTimeMillis();
out.add("");
out.add("===== MECHANISM (forced): allowRetry() called from a callback on ZK's event thread while policy=SUSPENDED and Curator is disconnected =====");
try (TestingServer server = new TestingServer(true)) {
AtomicReference<CuratorFramework> ref = new AtomicReference<>();
RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
RecordingPolicy rec = new RecordingPolicy(stock, ref::get, t0);
CuratorFramework zk = CuratorFrameworkFactory.builder().connectString(server.getConnectString())
.sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS).retryPolicy(rec).build();
ref.set(zk);
rec.bind(zk);
final java.util.concurrent.CountDownLatch suspended = new java.util.concurrent.CountDownLatch(1);
final long[] reconnectedAt = {0};
zk.getConnectionStateListenable().addListener((c, st) -> {
log("state -> " + st);
if (st == ConnectionState.SUSPENDED) {
suspended.countDown();
}
if (st == ConnectionState.RECONNECTED) {
reconnectedAt[0] = System.currentTimeMillis();
}
});
zk.start();
zk.blockUntilConnected(10, TimeUnit.SECONDS);
zk.create().forPath("/r");
final long[] cbStart = {0};
final long[] cbEnd = {0};
final boolean[] cbResult = {false};
final String[] cbThread = {null};
final int[] cbRc = {99};
log("stopping server");
server.stop();
suspended.await(10, TimeUnit.SECONDS);
Thread.sleep(300);
// Raw ZooKeeper async call: queued at the ZK level while disconnected, fails with CONNECTIONLOSS on a later
// failed connect attempt. Its callback runs on the event thread and calls allowRetry, as Curator's would.
org.apache.storm.shade.org.apache.zookeeper.ZooKeeper raw = zk.getZookeeperClient().getZooKeeper();
raw.getData("/r", false, (rc, path, ctx, data, stat) -> {
cbThread[0] = Thread.currentThread().getName();
cbRc[0] = rc;
cbStart[0] = System.currentTimeMillis();
log("raw callback on " + cbThread[0] + " rc=" + rc + " -> calling policy.allowRetry");
cbResult[0] = rec.allowRetry(0, 0L, (time, unit) -> unit.sleep(time));
cbEnd[0] = System.currentTimeMillis();
log("policy.allowRetry returned " + cbResult[0] + " after " + (cbEnd[0] - cbStart[0]) + " ms");
}, null);
long until = System.currentTimeMillis() + 8000;
while (cbStart[0] == 0 && System.currentTimeMillis() < until) {
Thread.sleep(20);
}
out.add("callback started: " + (cbStart[0] != 0) + " on thread " + cbThread[0] + " rc=" + cbRc[0]);
Thread.sleep(1500);
long restartedAt = System.currentTimeMillis();
server.restart();
log("server restarted (callback has been parked " + (restartedAt - cbStart[0]) + " ms)");
long end = System.currentTimeMillis() + 40_000;
while ((cbEnd[0] == 0 || reconnectedAt[0] == 0) && System.currentTimeMillis() < end) {
Thread.sleep(50);
}
Thread.sleep(300);
out.add("RESULT: allowRetry blocked " + (cbEnd[0] == 0 ? "NEVER RETURNED" : (cbEnd[0] - cbStart[0]) + " ms")
+ " and returned " + cbResult[0]);
out.add("RESULT: server was back " + (cbEnd[0] == 0 ? "?" : (cbEnd[0] - restartedAt) + " ms") + " before allowRetry returned");
out.add("RESULT: RECONNECTED notification " + (reconnectedAt[0] == 0 ? "NEVER" : (reconnectedAt[0] - restartedAt) + " ms after server restart"));
synchronized (rec.calls) {
rec.calls.forEach(c -> out.add(" " + c));
}
zk.close();
}
}
/** Steady traffic across the drop, to see whether the risky situation arises on its own. */
private void natural(boolean aware) throws Exception {
t0 = System.currentTimeMillis();
out.add("");
out.add("===== NATURAL: steady background traffic across a 4s outage | aware=" + aware + " =====");
try (TestingServer server = new TestingServer(true)) {
AtomicReference<CuratorFramework> ref = new AtomicReference<>();
RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
RecordingPolicy rec = aware ? new RecordingPolicy(stock, ref::get, t0) : null;
CuratorFramework zk = CuratorFrameworkFactory.builder().connectString(server.getConnectString())
.sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS).retryPolicy(aware ? rec : stock).build();
ref.set(zk);
if (aware) {
rec.bind(zk);
}
final long[] reconnectedAt = {0};
zk.getConnectionStateListenable().addListener((c, st) -> {
log("state -> " + st);
if (st == ConnectionState.RECONNECTED) {
reconnectedAt[0] = System.currentTimeMillis();
}
});
zk.start();
zk.blockUntilConnected(10, TimeUnit.SECONDS);
zk.create().forPath("/r");
AtomicInteger issued = new AtomicInteger();
AtomicInteger done = new AtomicInteger();
AtomicInteger ok = new AtomicInteger();
AtomicInteger lost = new AtomicInteger();
BackgroundCallback cb = (c, e) -> {
if (e.getResultCode() == KeeperException.Code.OK.intValue()) {
ok.incrementAndGet();
} else {
lost.incrementAndGet();
}
done.incrementAndGet();
};
java.util.concurrent.atomic.AtomicBoolean go = new java.util.concurrent.atomic.AtomicBoolean(true);
Thread issuer = new Thread(() -> {
while (go.get()) {
try {
zk.getData().inBackground(cb).forPath("/r");
issued.incrementAndGet();
} catch (Exception e) {
break;
}
java.util.concurrent.locks.LockSupport.parkNanos(20_000);
}
}, "issuer");
issuer.start();
Thread.sleep(300);
log("stopping server (traffic continues)");
server.stop();
Thread.sleep(600);
go.set(false);
issuer.join();
Thread.sleep(3400);
long restartedAt = System.currentTimeMillis();
server.restart();
log("server restarted");
long end = System.currentTimeMillis() + 40_000;
while ((reconnectedAt[0] == 0 || done.get() < issued.get()) && System.currentTimeMillis() < end) {
Thread.sleep(50);
}
Thread.sleep(300);
out.add("operations issued=" + issued + " completed=" + done + " ok=" + ok + " failed(retries exhausted)=" + lost);
out.add("RECONNECTED notification " + (reconnectedAt[0] == 0 ? "NEVER" : (reconnectedAt[0] - restartedAt) + " ms after server restart"));
if (aware) {
out.add("policy.allowRetry calls: total=" + rec.total + " onEventThread=" + rec.onEventThread
+ " onEventThreadWhileSuspendedOrLost=" + rec.onEventThreadWhileDown
+ " blockedOver500msOnEventThread=" + rec.blockedOnEventThread + " maxCallMs=" + rec.maxMs);
out.add(" OFF event thread: calls=" + rec.offEventThread + " whileSuspendedOrLost=" + rec.offEventThreadWhileDown
+ " longestWhileDownMs=" + rec.maxMsOffEventWhileDown + " threads=" + rec.offEventThreadNames);
synchronized (rec.calls) {
int n = 0;
for (String c : rec.calls) {
if (n++ < 8) {
out.add(" " + c);
}
}
if (rec.calls.size() > 8) {
out.add(" ... (" + rec.calls.size() + " notable calls total)");
}
}
}
zk.close();
}
}
/** Synchronous (foreground) operations across an outage that outlasts the (shortened) stock retry budget. */
private void foreground(boolean aware) throws Exception {
t0 = System.currentTimeMillis();
out.add("");
out.add("===== FOREGROUND: 16 threads doing synchronous getData across a 2s outage | aware=" + aware + " =====");
try (TestingServer server = new TestingServer(true)) {
AtomicReference<CuratorFramework> ref = new AtomicReference<>();
RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
RecordingPolicy rec = aware ? new RecordingPolicy(stock, ref::get, t0) : null;
CuratorFramework zk = CuratorFrameworkFactory.builder().connectString(server.getConnectString())
.sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS).retryPolicy(aware ? rec : stock).build();
ref.set(zk);
if (aware) {
rec.bind(zk);
}
final long[] reconnectedAt = {0};
zk.getConnectionStateListenable().addListener((c, st) -> {
log("state -> " + st);
if (st == ConnectionState.RECONNECTED) {
reconnectedAt[0] = System.currentTimeMillis();
}
});
zk.start();
zk.blockUntilConnected(10, TimeUnit.SECONDS);
zk.create().forPath("/r");
AtomicInteger ok = new AtomicInteger();
AtomicInteger failed = new AtomicInteger();
AtomicInteger maxMs = new AtomicInteger();
java.util.Map<String, Integer> failTypes = new java.util.concurrent.ConcurrentHashMap<>();
java.util.concurrent.atomic.AtomicBoolean go = new java.util.concurrent.atomic.AtomicBoolean(true);
List<Thread> threads = new ArrayList<>();
for (int i = 0; i < 16; i++) {
Thread t = new Thread(() -> {
while (go.get()) {
long a = System.currentTimeMillis();
try {
zk.getData().forPath("/r");
ok.incrementAndGet();
} catch (Exception e) {
failed.incrementAndGet();
failTypes.merge(e.getClass().getSimpleName(), 1, Integer::sum);
}
maxMs.accumulateAndGet((int) (System.currentTimeMillis() - a), Math::max);
java.util.concurrent.locks.LockSupport.parkNanos(1_000_000);
}
}, "fg-" + i);
t.start();
threads.add(t);
}
Thread.sleep(500);
log("stopping server");
server.stop();
Thread.sleep(2000);
long restartedAt = System.currentTimeMillis();
server.restart();
log("server restarted");
Thread.sleep(4000);
go.set(false);
for (Thread t : threads) {
t.join(15_000);
}
out.add("synchronous calls: ok=" + ok + " failed=" + failed + " " + failTypes + " longestCallMs=" + maxMs);
out.add("RECONNECTED notification " + (reconnectedAt[0] == 0 ? "NEVER" : (reconnectedAt[0] - restartedAt) + " ms after server restart"));
if (aware) {
out.add("policy.allowRetry calls: total=" + rec.total + " onEventThread=" + rec.onEventThread
+ " onEventThreadWhileSuspendedOrLost=" + rec.onEventThreadWhileDown
+ " blockedOver500msOnEventThread=" + rec.blockedOnEventThread + " maxCallMs=" + rec.maxMs);
}
zk.close();
}
}
@Test
public void repro() throws Exception {
try {
mechanism();
natural(false);
natural(true);
foreground(false);
foreground(true);
experiment("long outage (8s > 5s connection timeout)", false, 1, 8000, 500);
experiment("long outage (8s > 5s connection timeout)", true, 1, 8000, 500);
} finally {
try (PrintWriter w = new PrintWriter(OUT)) {
synchronized (out) {
out.forEach(w::println);
}
}
}
}
} |
Root cause: A transient ZooKeeper connection loss during the supervisor heartbeat cycle caused the process to die in 629ms via DefaultUncaughtExceptionHandler. The Curator RetryLoop blindly slept between retries, racing against the ZK client's SendThread reconnection instead of yielding to it.
Fix (two layers):
ConnectionAwareRetryPolicy (new): A RetryPolicy wrapper that tracks the ZK 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 (no magic numbers).
SupervisorHeartbeat safety net: Wrap the heartbeat body in try/catch so that a total ZK outage (beyond session timeout) skips the cycle instead of killing the process. A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (30s) allows 6 missed beats.
Wired up in CuratorUtils.newCurator() via AtomicReference late binding (the CuratorFramework doesn't exist until builder.build() returns).
What is the purpose of the change:
This fix prevents supervisor processes from dying due to transient ZooKeeper connection losses. During the incident, a single ZK ensemble member closing its socket caused the supervisor to exit in 629ms via
DefaultUncaughtExceptionHandler, killing all topology workers on that node. The fix absorbs these transient blips by making the Curator retry policy aware of the ZK connection state — whenSUSPENDEDorLOST, it yields to the ZK client's SendThread reconnection viablockUntilConnected()instead of blind sleep+retry. This ensures supervisors survive network/ZK blips that are already being handled by the ZK client's failover mechanism.How was the change tested:
The fix was deployed to a production cluster running 3 supervisors with 220+ active workers. During a subsequent network outage that caused multiple transient ZK connection losses, the supervisors survived without dying (previously they would have exited via
DefaultUncaughtExceptionHandler). Log verification confirmed theConnectionAwareRetryPolicycorrectly detectedSUSPENDEDstates, waited for reconnection viablockUntilConnected(), and retried operations on the new connection. No supervisors died, no topology workers were killed due to supervisor crashes.Closes #9098
EDIT 2026-10-04: attach crash log
connectionloss-crash-anonymized.log