reiabreu commented on code in PR #9153:
URL: https://github.com/apache/storm/pull/9153#discussion_r4172977344


##########
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.
+ *
+ * <p>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.
+ *
+ * <p>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.
+ *
+ * <p>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<CuratorFramework> zkSupplier;
+    private final int sessionTimeoutMs;
+    private final AtomicReference<ConnectionState> 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<CuratorFramework> 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;
+            }

Review 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.
   
   ```suggestion
               // 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;
               }
   ```



-- 
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]

Reply via email to