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

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


The following commit(s) were added to refs/heads/master by this push:
     new e3f06eac3ec fix: per-segment timeout accounting bug (#20290)
e3f06eac3ec is described below

commit e3f06eac3ec902937f72caa31871cf97d5a1df12
Author: jtuglu1 <[email protected]>
AuthorDate: Tue Sep 15 08:27:46 2026 +0300

    fix: per-segment timeout accounting bug (#20290)
---
 .../druid/query/ForwardingQueryProcessingPool.java |  49 ++++-
 .../query/ForwardingQueryProcessingPoolTest.java   | 233 +++++++++++++++++++++
 2 files changed, 274 insertions(+), 8 deletions(-)

diff --git 
a/processing/src/main/java/org/apache/druid/query/ForwardingQueryProcessingPool.java
 
b/processing/src/main/java/org/apache/druid/query/ForwardingQueryProcessingPool.java
index 532cdd85f4d..f98897114b0 100644
--- 
a/processing/src/main/java/org/apache/druid/query/ForwardingQueryProcessingPool.java
+++ 
b/processing/src/main/java/org/apache/druid/query/ForwardingQueryProcessingPool.java
@@ -25,6 +25,7 @@ import com.google.common.util.concurrent.Futures;
 import com.google.common.util.concurrent.ListenableFuture;
 import com.google.common.util.concurrent.ListeningExecutorService;
 import com.google.common.util.concurrent.MoreExecutors;
+import com.google.common.util.concurrent.SettableFuture;
 
 import javax.annotation.Nullable;
 import java.util.concurrent.ExecutorService;
@@ -63,6 +64,11 @@ public class ForwardingQueryProcessingPool extends 
ForwardingListeningExecutorSe
     return delegate().submit(task);
   }
 
+  /**
+   * The timeout only starts counting once a processing thread has actually 
picked the task off the pool's queue,
+   * not when it is submitted. Otherwise a task that sits in the queue behind 
other segments could exhaust its
+   * per-segment timeout without ever having been given a chance to run.
+   */
   @Override
   public <T, V> ListenableFuture<T> submitRunnerTask(
       PrioritizedQueryRunnerCallable<T, V> task,
@@ -70,15 +76,42 @@ public class ForwardingQueryProcessingPool extends 
ForwardingListeningExecutorSe
       TimeUnit unit
   )
   {
-    if (timeoutService != null) {
-      return Futures.withTimeout(
-          delegate().submit(task),
-          timeout,
-          unit,
-          timeoutService
-      );
+    if (timeoutService == null) {
+      return submitRunnerTask(task);
     }
-    return submitRunnerTask(task);
+
+    final SettableFuture<Void> started = SettableFuture.create();
+    final ListenableFuture<T> execFuture = submitRunnerTask(
+        new AbstractPrioritizedQueryRunnerCallable<T, V>(task.getPriority(), 
task.getRunner())
+        {
+          @Override
+          public T call() throws Exception
+          {
+            started.set(null);
+            return task.call();
+          }
+        }
+    );
+    // If the task never gets to run (cancelled or rejected while queued), 
unblock the transform below so that the
+    // returned future completes with the underlying outcome instead of 
hanging forever.
+    execFuture.addListener(() -> started.set(null), 
MoreExecutors.directExecutor());
+
+    final ListenableFuture<T> timedFuture = Futures.transformAsync(
+        started,
+        ignored -> Futures.withTimeout(execFuture, timeout, unit, 
timeoutService),
+        MoreExecutors.directExecutor()
+    );
+    // Cancelling the returned future while the task is still queued must 
cancel the queued task as well, which the
+    // transform cannot do on its own since it is waiting on 'started' rather 
than on the task itself.
+    timedFuture.addListener(
+        () -> {
+          if (timedFuture.isCancelled()) {
+            execFuture.cancel(true);
+          }
+        },
+        MoreExecutors.directExecutor()
+    );
+    return timedFuture;
   }
 
   @Override
diff --git 
a/processing/src/test/java/org/apache/druid/query/ForwardingQueryProcessingPoolTest.java
 
b/processing/src/test/java/org/apache/druid/query/ForwardingQueryProcessingPoolTest.java
new file mode 100644
index 00000000000..313d5368bd8
--- /dev/null
+++ 
b/processing/src/test/java/org/apache/druid/query/ForwardingQueryProcessingPoolTest.java
@@ -0,0 +1,233 @@
+/*
+ * 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.druid.query;
+
+import com.google.common.util.concurrent.ListenableFuture;
+import org.apache.druid.java.util.common.concurrent.Execs;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+public class ForwardingQueryProcessingPoolTest
+{
+  private ExecutorService exec;
+  private ScheduledExecutorService timeoutService;
+  private ForwardingQueryProcessingPool pool;
+
+  @BeforeEach
+  public void setUp()
+  {
+    // Single threaded, so that a second task submitted has to wait in the 
queue for the first one to finish.
+    exec = Execs.singleThreaded("ForwardingQueryProcessingPoolTest-%d");
+    timeoutService = 
Execs.scheduledSingleThreaded("ForwardingQueryProcessingPoolTest-timeout-%d");
+    pool = new ForwardingQueryProcessingPool(exec, timeoutService);
+  }
+
+  @AfterEach
+  public void tearDown()
+  {
+    pool.shutdownNow();
+    exec.shutdownNow();
+    timeoutService.shutdownNow();
+  }
+
+  /**
+   * The timeout of a task must not run down while the task is still waiting 
in the pool's queue.
+   */
+  @Test
+  @Timeout(30)
+  public void testTimeoutStartsWhenTaskIsPickedUpAndNotWhenSubmitted() throws 
Exception
+  {
+    final CountDownLatch blockerStarted = new CountDownLatch(1);
+    final CountDownLatch releaseBlocker = new CountDownLatch(1);
+
+    final ListenableFuture<Integer> blockerFuture = 
pool.submitRunnerTask(callable(() -> {
+      blockerStarted.countDown();
+      releaseBlocker.await();
+      return 1;
+    }));
+    blockerStarted.await();
+
+    // Queued behind the blocker, with a timeout much shorter than the time it 
spends waiting.
+    final ListenableFuture<Integer> queuedFuture = pool.submitRunnerTask(
+        callable(() -> 2),
+        500,
+        TimeUnit.MILLISECONDS
+    );
+
+    Assertions.assertFalse(queuedFuture.isDone());
+    Thread.sleep(1500);
+    Assertions.assertFalse(queuedFuture.isDone(), "Queued task should not time 
out before it starts running");
+
+    releaseBlocker.countDown();
+    Assertions.assertEquals(1, blockerFuture.get().intValue());
+    Assertions.assertEquals(2, queuedFuture.get().intValue());
+  }
+
+  @Test
+  @Timeout(30)
+  public void testTimeoutFiresOnceTheTaskIsRunningForTooLong()
+  {
+    final ListenableFuture<Integer> future = pool.submitRunnerTask(
+        callable(() -> {
+          Thread.sleep(60_000L);
+          return 1;
+        }),
+        100,
+        TimeUnit.MILLISECONDS
+    );
+
+    final ExecutionException e = 
Assertions.assertThrows(ExecutionException.class, future::get);
+    Assertions.assertInstanceOf(TimeoutException.class, e.getCause());
+  }
+
+  @Test
+  @Timeout(30)
+  public void testTaskCompletesNormallyWithinTimeout() throws Exception
+  {
+    Assertions.assertEquals(
+        1,
+        pool.submitRunnerTask(callable(() -> 1), 30_000, 
TimeUnit.MILLISECONDS).get().intValue()
+    );
+  }
+
+  /**
+   * Cancelling the returned future must cancel the underlying task even while 
it is still sitting in the queue.
+   */
+  @Test
+  @Timeout(30)
+  public void testCancelWhileQueued() throws Exception
+  {
+    final CountDownLatch blockerStarted = new CountDownLatch(1);
+    final CountDownLatch releaseBlocker = new CountDownLatch(1);
+    final AtomicBoolean queuedTaskRan = new AtomicBoolean(false);
+
+    final ListenableFuture<Integer> blockerFuture = 
pool.submitRunnerTask(callable(() -> {
+      blockerStarted.countDown();
+      releaseBlocker.await();
+      return 1;
+    }));
+    blockerStarted.await();
+
+    final ListenableFuture<Integer> queuedFuture = pool.submitRunnerTask(
+        callable(() -> {
+          queuedTaskRan.set(true);
+          return 2;
+        }),
+        30_000,
+        TimeUnit.MILLISECONDS
+    );
+    Assertions.assertTrue(queuedFuture.cancel(true));
+
+    releaseBlocker.countDown();
+    Assertions.assertEquals(1, blockerFuture.get().intValue());
+    // Give the pool a chance to (incorrectly) run the cancelled task.
+    Thread.sleep(500);
+    Assertions.assertFalse(queuedTaskRan.get(), "Cancelled task should not 
have been run");
+  }
+
+  /**
+   * A cancelled task should not leave the returned future hanging, since it 
never gets to start.
+   */
+  @Test
+  @Timeout(30)
+  public void testCancelWhileQueuedCompletesReturnedFuture() throws Exception
+  {
+    final CountDownLatch blockerStarted = new CountDownLatch(1);
+    final CountDownLatch releaseBlocker = new CountDownLatch(1);
+
+    final ListenableFuture<Integer> blockerFuture = 
pool.submitRunnerTask(callable(() -> {
+      blockerStarted.countDown();
+      releaseBlocker.await();
+      return 1;
+    }));
+    blockerStarted.await();
+
+    final ListenableFuture<Integer> queuedFuture = pool.submitRunnerTask(
+        callable(() -> 2),
+        30_000,
+        TimeUnit.MILLISECONDS
+    );
+    queuedFuture.cancel(true);
+    releaseBlocker.countDown();
+
+    Assertions.assertEquals(1, blockerFuture.get().intValue());
+    Assertions.assertTrue(queuedFuture.isDone());
+    Assertions.assertTrue(queuedFuture.isCancelled());
+  }
+
+  @Test
+  @Timeout(30)
+  public void testPriorityAndRunnerArePreserved()
+  {
+    final QueryRunner<Object> runner = (queryPlus, responseContext) -> null;
+    final PrioritizedQueryRunnerCallable<Integer, Object> task =
+        new AbstractPrioritizedQueryRunnerCallable<>(7, runner)
+        {
+          @Override
+          public Integer call()
+          {
+            return 1;
+          }
+        };
+
+    final PrioritizedQueryRunnerCallable<?, ?>[] submitted = new 
PrioritizedQueryRunnerCallable<?, ?>[1];
+    final ForwardingQueryProcessingPool capturingPool = new 
ForwardingQueryProcessingPool(exec, timeoutService)
+    {
+      @Override
+      public <T, V> ListenableFuture<T> 
submitRunnerTask(PrioritizedQueryRunnerCallable<T, V> task)
+      {
+        submitted[0] = task;
+        return super.submitRunnerTask(task);
+      }
+    };
+
+    capturingPool.submitRunnerTask(task, 30_000, TimeUnit.MILLISECONDS);
+    Assertions.assertEquals(7, submitted[0].getPriority());
+    Assertions.assertSame(runner, submitted[0].getRunner());
+  }
+
+  private static PrioritizedQueryRunnerCallable<Integer, Object> 
callable(ThrowingSupplier supplier)
+  {
+    return new AbstractPrioritizedQueryRunnerCallable<>(0, null)
+    {
+      @Override
+      public Integer call() throws Exception
+      {
+        return supplier.get();
+      }
+    };
+  }
+
+  private interface ThrowingSupplier
+  {
+    Integer get() throws Exception;
+  }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to