This is an automated email from the ASF dual-hosted git repository.

rzo1 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/storm.git


The following commit(s) were added to refs/heads/master by this push:
     new 4af262d17 Remove DRPC per-function request queues once they are empty 
(#8985)
4af262d17 is described below

commit 4af262d177eaf5053ca9d00efd46f79e20490244
Author: Richard Zowalla <[email protected]>
AuthorDate: Sat Aug 22 20:44:47 2026 +0200

    Remove DRPC per-function request queues once they are empty (#8985)
---
 .../java/org/apache/storm/daemon/drpc/DRPC.java    |  48 +++++++---
 .../org/apache/storm/daemon/drpc/DRPCTest.java     | 105 +++++++++++++++++++++
 2 files changed, 140 insertions(+), 13 deletions(-)

diff --git a/storm-server/src/main/java/org/apache/storm/daemon/drpc/DRPC.java 
b/storm-server/src/main/java/org/apache/storm/daemon/drpc/DRPC.java
index 23183f08e..df9ed364f 100644
--- a/storm-server/src/main/java/org/apache/storm/daemon/drpc/DRPC.java
+++ b/storm-server/src/main/java/org/apache/storm/daemon/drpc/DRPC.java
@@ -28,6 +28,7 @@ import java.util.TimerTask;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentLinkedQueue;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
 import org.apache.storm.DaemonConfig;
 import org.apache.storm.daemon.StormCommon;
 import org.apache.storm.generated.AuthorizationException;
@@ -145,8 +146,15 @@ public class DRPC implements AutoCloseable {
 
     private void cleanup(String id) {
         OutstandingRequest req = requests.remove(id);
-        if (req != null && !req.wasFetched()) {
-            queues.get(req.getFunction()).remove(req);
+        if (req != null) {
+            queues.computeIfPresent(req.getFunction(), (function, queue) -> {
+                if (!req.wasFetched()) {
+                    queue.remove(req);
+                }
+                //Drop the queue itself once nothing is waiting in it, 
otherwise the map keeps an
+                // entry for every function name a client has ever asked about.
+                return queue.isEmpty() ? null : queue;
+            });
         }
     }
 
@@ -165,16 +173,15 @@ public class DRPC implements AutoCloseable {
         return String.valueOf(ctr.incrementAndGet());
     }
 
-    private ConcurrentLinkedQueue<OutstandingRequest> getQueue(String 
function) {
+    private static void checkFunctionName(String function) {
         if (function == null) {
             throw new IllegalArgumentException("The function for a request 
cannot be null");
         }
-        ConcurrentLinkedQueue<OutstandingRequest> queue = queues.get(function);
-        if (queue == null) {
-            queues.putIfAbsent(function, new ConcurrentLinkedQueue<>());
-            queue = queues.get(function);
-        }
-        return queue;
+    }
+
+    @VisibleForTesting
+    int getNumTrackedFunctions() {
+        return queues.size();
     }
 
     public void returnResult(String id, String result) throws 
AuthorizationException {
@@ -190,8 +197,17 @@ public class DRPC implements AutoCloseable {
     public DRPCRequest fetchRequest(String functionName) throws 
AuthorizationException {
         meterFetchRequestCalls.mark();
         checkAuthorizationNoLog("fetchRequest", functionName);
-        ConcurrentLinkedQueue<OutstandingRequest> q = getQueue(functionName);
-        OutstandingRequest req = q.poll();
+        checkFunctionName(functionName);
+        //Never create a queue here.  A function name comes from the client, 
so a queue that no one
+        // ever puts a request into would stay in the map forever.  Poll and 
drop an emptied queue
+        // under the same lock execute() adds under, so a request can never be 
left in a queue that
+        // was just removed from the map.
+        AtomicReference<OutstandingRequest> polled = new AtomicReference<>();
+        queues.computeIfPresent(functionName, (function, queue) -> {
+            polled.set(queue.poll());
+            return queue.isEmpty() ? null : queue;
+        });
+        OutstandingRequest req = polled.get();
         if (req != null) {
             //Only log accesses that fetched something
             logAccess("fetchRequest", functionName);
@@ -219,12 +235,18 @@ public class DRPC implements AutoCloseable {
         AuthorizationException {
         meterExecuteCalls.mark();
         checkAuthorization("execute", functionName);
+        checkFunctionName(functionName);
         String id = nextId();
         LOG.debug("Execute {} {}", functionName, funcArgs);
         T req = factory.mkRequest(functionName, new DRPCRequest(funcArgs, id));
         requests.put(id, req);
-        ConcurrentLinkedQueue<OutstandingRequest> q = getQueue(functionName);
-        q.add(req);
+        queues.compute(functionName, (function, queue) -> {
+            if (queue == null) {
+                queue = new ConcurrentLinkedQueue<>();
+            }
+            queue.add(req);
+            return queue;
+        });
         return req;
     }
 
diff --git 
a/storm-server/src/test/java/org/apache/storm/daemon/drpc/DRPCTest.java 
b/storm-server/src/test/java/org/apache/storm/daemon/drpc/DRPCTest.java
index 210aa7246..a9fd7f1d1 100644
--- a/storm-server/src/test/java/org/apache/storm/daemon/drpc/DRPCTest.java
+++ b/storm-server/src/test/java/org/apache/storm/daemon/drpc/DRPCTest.java
@@ -18,9 +18,13 @@
 
 package org.apache.storm.daemon.drpc;
 
+import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
@@ -140,6 +144,107 @@ public class DRPCTest {
         }
     }
 
+    @Test
+    public void testQueuesAreRemovedWhenEmpty() throws Exception {
+        try (DRPC server = new DRPC(new StormMetricsRegistry(), null, 1000)) {
+            //Fetching for a function nothing was ever submitted for must not 
leave state behind
+            DRPCRequest nothing = server.fetchRequest("never-registered");
+            assertNotNull(nothing);
+            assertEquals("", nothing.get_request_id());
+            assertEquals(0, server.getNumTrackedFunctions());
+
+            //A registered function is still served repeatedly, and is not 
left behind once idle
+            for (int i = 0; i < 3; i++) {
+                Future<String> found = exec.submit(() -> 
server.executeBlocking("testing", "test"));
+                DRPCRequest request = getNextAvailableRequest(server, 
"testing");
+                assertNotNull(request);
+                server.returnResult(request.get_request_id(), "tested");
+                assertEquals("tested", found.get(10, TimeUnit.MILLISECONDS));
+            }
+            assertEquals(0, server.getNumTrackedFunctions());
+
+            //Nor is a function whose only request timed out.  The timer 
thread fails the request
+            // before it drops the queue, so the caller can return first; wait 
for the drop instead
+            // of racing it, with a hard timeout so a real leak still fails 
the test.
+            try {
+                server.executeBlocking("timing-out", "test");
+                fail("Should have timed out....");
+            } catch (DRPCExecutionException e) {
+                assertEquals(DRPCExceptionType.SERVER_TIMEOUT, e.get_type());
+            }
+            Awaitility.await("DRPC queue for timing-out to be dropped")
+                .atMost(5, TimeUnit.SECONDS)
+                .pollInterval(1, TimeUnit.MILLISECONDS)
+                .until(() -> server.getNumTrackedFunctions() == 0);
+        }
+    }
+
+    @Test
+    public void testConcurrentExecuteAndFetchLosesNoRequests() throws 
Exception {
+        //A bounded pool of 16 threads is what keeps this cheap: 
executeBlocking() parks its caller,
+        // so an unbounded pool would need one live thread per outstanding 
request.  The request
+        // count costs no threads at all, and is what gives the stress test 
its power.  Measured
+        // against a fetchRequest() whose poll/remove escapes the per-function 
compute lock, the
+        // lost request was caught 1 run in 25 at 200 requests, 3 in 10 at 
2000 and 9 in 10 at 5000,
+        // while a correct server still serves all 5000 in about a second with 
every core busy.
+        final int numRequests = 5000;
+        final int numThreads = 16;
+        final long deadlineMs = 30_000;
+        //A timeout far beyond the test deadline, so the cleanup timer never 
reaps a live request.
+        try (DRPC server = new DRPC(new StormMetricsRegistry(), null, 
300_000)) {
+            ExecutorService submitters = 
Executors.newFixedThreadPool(numThreads);
+            try {
+                List<Future<String>> futures = new ArrayList<>(numRequests);
+                for (int i = 0; i < numRequests; i++) {
+                    final String args = "test-" + i;
+                    futures.add(submitters.submit(() -> 
server.executeBlocking("testing", args)));
+                }
+
+                Set<String> servedIds = new HashSet<>();
+                long deadline = Time.currentTimeMillis() + deadlineMs;
+                int emptyFetches = 0;
+                while (servedIds.size() < numRequests) {
+                    if (Time.currentTimeMillis() > deadline) {
+                        fail("Only served " + servedIds.size() + " of " + 
numRequests
+                             + " requests within " + deadlineMs + "ms, a 
request was lost");
+                    }
+                    DRPCRequest req = server.fetchRequest("testing");
+                    assertNotNull(req);
+                    String id = req.get_request_id();
+                    if (id.isEmpty()) {
+                        //Nothing to serve right now.  Spin at first, so 
fetches keep interleaving
+                        // tightly with the submitting threads, and only back 
off if this goes on
+                        // for a long time (a regression, which the deadline 
above then fails).
+                        if (++emptyFetches > 10_000) {
+                            TimeUnit.MILLISECONDS.sleep(1);
+                        } else {
+                            Thread.onSpinWait();
+                        }
+                        continue;
+                    }
+                    emptyFetches = 0;
+                    assertTrue(servedIds.add(id), "Request " + id + " was 
fetched more than once");
+                    server.returnResult(id, "tested-" + id);
+                }
+
+                Set<String> results = new HashSet<>();
+                for (Future<String> f : futures) {
+                    long left = deadline - Time.currentTimeMillis();
+                    assertTrue(left > 0, "Ran out of time waiting for the 
blocked callers");
+                    assertTrue(results.add(f.get(left, 
TimeUnit.MILLISECONDS)), "Duplicate result returned");
+                }
+                assertEquals(numRequests, results.size());
+                for (String id : servedIds) {
+                    assertTrue(results.contains("tested-" + id), "No caller 
got the result for " + id);
+                }
+                //Nothing is waiting any more, so no per-function queue may be 
left behind
+                assertEquals(0, server.getNumTrackedFunctions());
+            } finally {
+                submitters.shutdownNow();
+            }
+        }
+    }
+
     @Test
     public void testDeny() {
         try (DRPC server = new DRPC(new StormMetricsRegistry(), new 
DenyAuthorizer(), 100)) {

Reply via email to