diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java new file mode 100644 index 00000000000..eb041bdefde --- /dev/null +++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java @@ -0,0 +1,146 @@ +/* + * 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.ThreadLocalRandom; +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.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.state.ConnectionState; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * A {@link RetryPolicy} wrapper that makes the Curator retry loop aware of the + * ZooKeeper connection state. + * + *

When the connection is {@link ConnectionState#SUSPENDED} or + * {@link ConnectionState#LOST}, instead of blindly sleeping and retrying (which + * races against the ZK client's {@code SendThread} reconnection), this policy + * calls {@link CuratorFramework#blockUntilConnected} to yield to the + * {@code SendThread} and wait for it to failover to another ensemble member. + * + *

Once the connection is re-established, the retry loop immediately retries + * the operation on the new connection. If the connection cannot be + * re-established within the session timeout, the retry is abandoned. + * + *

For all other connection states, this policy delegates to the wrapped + * delegate policy (typically a {@link StormBoundedExponentialBackoffRetry}). + */ +public class ConnectionAwareRetryPolicy implements RetryPolicy { + + private static final Logger LOG = LoggerFactory.getLogger(ConnectionAwareRetryPolicy.class); + + /** + * Maximum jitter (in ms) applied after a successful reconnection to avoid + * thundering-herd when many suspended clients retry in lockstep. + */ + private static final long RECONNECT_JITTER_MS = 100; + + private final RetryPolicy delegate; + private final Supplier zkSupplier; + private final int sessionTimeoutMs; + private final AtomicReference connectionState = + new AtomicReference<>(ConnectionState.CONNECTED); + + /** + * @param delegate the underlying retry policy to delegate to for normal (connected) retries + * @param zkSupplier supplier for the {@link CuratorFramework}, used for late binding since + * 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} + */ + public ConnectionAwareRetryPolicy(RetryPolicy delegate, + Supplier zkSupplier, + int sessionTimeoutMs) { + this.delegate = delegate; + this.zkSupplier = zkSupplier; + this.sessionTimeoutMs = sessionTimeoutMs; + } + + /** + * Register a {@link org.apache.storm.shade.org.apache.curator.framework.state.ConnectionStateListener} + * on the given framework to track the current connection state. + * Must be called after the {@link CuratorFramework} has been built. + * + * @param zk the built CuratorFramework + */ + public void bind(CuratorFramework zk) { + zk.getConnectionStateListenable().addListener((client, newState) -> { + ConnectionState prev = connectionState.getAndSet(newState); + if (prev != newState) { + LOG.debug("ZK connection state changed: {} -> {}", prev, newState); + } + }); + } + + @VisibleForTesting + RetryPolicy getDelegate() { + return delegate; + } + + @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)) { + return false; + } + + CuratorFramework zk = zkSupplier.get(); + if (zk == null) { + // Framework not yet available; delegate already approved the retry + return true; + } + + LOG.info("ZK connection is {} on retry {}, waiting for reconnection (timeout {}ms)", + state, retryCount, sessionTimeoutMs); + try { + boolean reconnected = zk.blockUntilConnected(sessionTimeoutMs, TimeUnit.MILLISECONDS); + if (reconnected) { + LOG.info("ZK connection re-established (state: {}), retrying operation", + connectionState.get()); + // Jitter to avoid thundering-herd on ensemble recovery. + long jitterMs = ThreadLocalRandom.current().nextLong(RECONNECT_JITTER_MS); + if (jitterMs > 0) { + sleepSleeper.sleepFor(jitterMs, TimeUnit.MILLISECONDS); + } + return true; + } + LOG.warn("ZK connection not re-established within {}ms, abandoning retry", + sessionTimeoutMs); + return false; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + LOG.warn("Interrupted while waiting for ZK reconnection", e); + return false; + } + } + + // Connection is healthy (CONNECTED, RECONNECTED, or READ_ONLY) — + // delegate to the existing backoff policy for normal retry behaviour. + return delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper); + } +} 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 d5e6f3f3e0d..d2365bd1ac0 100644 --- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java +++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; import javax.naming.ConfigurationException; import org.apache.storm.Config; import org.apache.storm.shade.org.apache.commons.lang3.StringUtils; @@ -58,6 +59,19 @@ public static CuratorFramework newCurator(Map conf, List CuratorFrameworkFactory.Builder builder = CuratorFrameworkFactory.builder(); setupBuilder(builder, zkStr, conf, auth); + + // Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST. + AtomicReference zkRef = new AtomicReference<>(); + int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT)); + 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))), + zkRef::get, + sessionTimeoutMs); + builder.retryPolicy(connectionAwarePolicy); + if (defaultAcl != null) { builder.aclProvider(new ACLProvider() { @Override @@ -72,11 +86,15 @@ public List getAclForPath(String s) { }); } - return builder.build(); + CuratorFramework framework = builder.build(); + zkRef.set(framework); + connectionAwarePolicy.bind(framework); + return framework; } protected static void setupBuilder(CuratorFrameworkFactory.Builder builder, final String zkStr, Map conf, ZookeeperAuthInfo auth) { + // Default retry policy; overridden by newCurator() with ConnectionAwareRetryPolicy. builder.connectString(zkStr); builder .connectionTimeoutMs(ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_CONNECTION_TIMEOUT))) diff --git a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java new file mode 100644 index 00000000000..38236fffa5c --- /dev/null +++ b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java @@ -0,0 +1,174 @@ +/* + * 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.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; + + private final RetryPolicy delegate = mock(RetryPolicy.class); + private final RetrySleeper sleeper = mock(RetrySleeper.class); + private final CuratorFramework zk = mock(CuratorFramework.class); + private ConnectionStateListener listener; + + @SuppressWarnings("unchecked") + private ConnectionAwareRetryPolicy build(Supplier zkSupplier) { + Listenable listenable = mock(Listenable.class); + when(zk.getConnectionStateListenable()).thenReturn(listenable); + + ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS); + policy.bind(zk); + + ArgumentCaptor 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(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); + + assertTrue(result, "should retry once the SendThread has reconnected"); + verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS); + } + + @Test + public void lostBlocksUntilConnectedThenRetriesImmediately() 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); + + assertTrue(result, "LOST should also wait for reconnection then retry"); + verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS); + } + + @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"); + } + + @Test + public void suspendedHonoursDelegateRetryBudget() 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); + + assertFalse(result, "should abandon when the delegate's retry budget is exhausted"); + verify(zk, never()).blockUntilConnected(anyInt(), 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); + + 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); + } +} diff --git a/storm-client/test/jvm/org/apache/storm/utils/CuratorUtilsTest.java b/storm-client/test/jvm/org/apache/storm/utils/CuratorUtilsTest.java index 68f730fc6fe..c2745cb26ee 100644 --- a/storm-client/test/jvm/org/apache/storm/utils/CuratorUtilsTest.java +++ b/storm-client/test/jvm/org/apache/storm/utils/CuratorUtilsTest.java @@ -35,6 +35,7 @@ import org.apache.storm.shade.org.apache.zookeeper.ZooKeeper; import org.apache.storm.shade.org.apache.zookeeper.client.ZKClientConfig; import org.apache.storm.shade.org.apache.zookeeper.common.ClientX509Util; +import org.apache.storm.shade.org.apache.curator.RetryPolicy; import org.junit.jupiter.api.Test; import org.apache.curator.test.TestingServer; import org.slf4j.Logger; @@ -43,6 +44,7 @@ import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; public class CuratorUtilsTest { private static final Logger LOG = LoggerFactory.getLogger(CuratorUtilsTest.class); @@ -71,8 +73,11 @@ public void newCuratorUsesExponentialBackoffTest() { CuratorFramework curator = CuratorUtils.newCurator(config, Collections.singletonList("bogus_server"), 42, "", DaemonType.WORKER.getDefaultZkAcls(config)); + RetryPolicy retryPolicy = curator.getZookeeperClient().getRetryPolicy(); + assertTrue(retryPolicy instanceof ConnectionAwareRetryPolicy, + "newCurator should wrap the retry policy in a ConnectionAwareRetryPolicy"); StormBoundedExponentialBackoffRetry policy = - (StormBoundedExponentialBackoffRetry) curator.getZookeeperClient().getRetryPolicy(); + (StormBoundedExponentialBackoffRetry) ((ConnectionAwareRetryPolicy) retryPolicy).getDelegate(); assertEquals(policy.getBaseSleepTimeMs(), expectedInterval); assertEquals(policy.getN(), expectedRetries); assertEquals(policy.getSleepTimeMs(10, 0), expectedCeiling); diff --git a/storm-server/src/main/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeat.java b/storm-server/src/main/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeat.java index b613f36ee9a..3ee2c7c7ab7 100644 --- a/storm-server/src/main/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeat.java +++ b/storm-server/src/main/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeat.java @@ -162,7 +162,13 @@ public void run() { Map validatedNumaMap = SupervisorUtils.getNumaMap(conf); Map supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap); for (Map.Entry supervisorInfoEntry : supervisorInfoList.entrySet()) { - stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue()); + try { + stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue()); + } catch (Exception e) { + // Liveness is tracked via ephemeral ZK nodes, not this write. Safe to skip. + LOG.warn("Supervisor ZK heartbeat failed for {}, will retry next cycle", + supervisorInfoEntry.getKey(), e); + } } } } diff --git a/storm-server/src/test/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeatTest.java b/storm-server/src/test/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeatTest.java new file mode 100644 index 00000000000..bc25f2a03b6 --- /dev/null +++ b/storm-server/src/test/java/org/apache/storm/daemon/supervisor/timer/SupervisorHeartbeatTest.java @@ -0,0 +1,87 @@ +/* + * 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"; + + 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 baseConf() { + return new HashMap<>(Utils.readStormConfig()); + } + + @Test + public void transientHeartbeatFailureIsSwallowedInsteadOfKillingTheProcess() { + IStormClusterState clusterState = mock(IStormClusterState.class); + doThrow(new RuntimeException("simulated transient ZK failure")) + .when(clusterState).supervisorHeartbeat(anyString(), any(SupervisorInfo.class)); + + Supervisor supervisor = mockSupervisor(clusterState); + SupervisorHeartbeat heartbeat = new SupervisorHeartbeat(baseConf(), supervisor); + + assertDoesNotThrow(heartbeat::run); + 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)); + } +}