reiabreu commented on PR #9153:
URL: https://github.com/apache/storm/pull/9153#issuecomment-5973530933
*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 `ConnectionAwareRetryPolicy`
as written can block ZooKeeper's event thread, which delays the "reconnected"
event. A fix for this PR is below. With it, the policy behaved like stock
retries in the same tests.
**Problem.** For background (`inBackground`) operations, the kind Nimbus's
`LeaderLatch` uses, Curator calls `allowRetry` on ZooKeeper's event thread.
`blockUntilConnected` waits for a flag that only an event on that same thread
can set, so the wait can't end early. Claude ran the policy against a real
ZooKeeper (Curator's `TestingServer`, Storm's shaded client), killing and
restarting the server under load, next to a stock-backoff control. Steady
background reads across a 4 s outage:
| | stock backoff | as written | with the fix |
|---|---|---|---|
| operations failed | 0 | 2 | 0 |
| "reconnected" notice after the server was back | 0.5-1.4 s | **16.3 s** |
0.47 s |
| event-thread blocking | none | 2 blocks, about 18 s | none |
| session | up | went `LOST` | up |
A forced version is deterministic: `allowRetry` on the event thread blocked
10,014 ms and returned `false` although the server had been back for 8.5 s.
With the fix it returns in 107 ms.
**Fix for this PR.** Wait for the reconnect only on foreground retries and
delegate background retries to the wrapped policy. Curator passes
`RetryLoop.getDefaultRetrySleeper()` for foreground retries and the operation
itself for background ones (Curator 5.9.0 source). If that check ever stops
matching, the policy falls back to stock behavior.
```java
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 `maxRetries` in
`CuratorUtils`, and updates `ConnectionAwareRetryPolicyTest` to 10 cases,
including two that assert background retries never block. These tests and
`CuratorUtilsTest` pass locally.
**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.**
- Claude could not reproduce the original crash in a test. Stock retries
survived a 2 s outage with 16 threads of synchronous reads, so the foreground
benefit rests on the production report and the unit tests.
- Traffic was synthetic (about 20,000 operations per second), one run per
cell. Nimbus HA and session expiry are untested.
**Re @GGraziadei's question** (was the policy consulted?). The posted trace
can't tell. `RetryLoop.java:88` is `proc.call()`, and a refused retry rethrows
the same exception object. The supervisor's effective
`storm.zookeeper.retry.times` / `.interval`, or Curator's DEBUG logs from
`RetryLoopImpl`, would. The log also shows SUSPENDED periods longer than the 20
s session timeout (36 s, 20 s, 36 s), where waiting can't help; the `catch` is
what keeps the process alive. If maintainers prefer this PR to stay with the
`catch`, the fix above can go in a follow-up.
<details><summary>Patch (on top of ab67921da)</summary>
```diff
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());
}
}
```
</details>
<details><summary>Failover test program (throwaway JUnit test, not for
merging)</summary>
```java
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);
}
}
}
}
}
```
</details>
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]