From 314aed0bc1c84560409fd9a9a3565c4b34b05f77 Mon Sep 17 00:00:00 2001 From: JackieTien97 Date: Thu, 11 Jun 2026 18:10:38 +0800 Subject: [PATCH 1/2] Fix driver scheduler reservation leak --- .../execution/schedule/DriverScheduler.java | 11 +-- .../queue/IndexedBlockingReserveQueue.java | 32 ++++++++- .../MultilevelPriorityQueue.java | 10 +++ .../execution/schedule/task/DriverTask.java | 13 ++++ .../schedule/DefaultDriverSchedulerTest.java | 69 +++++++++++++++++++ 5 files changed, 128 insertions(+), 7 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java index 183e51b7e1ef1..03110a8f1dd83 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java @@ -336,16 +336,18 @@ private void clearDriverTask(DriverTask task) { return; case READY: task.setStatus(DriverTaskStatus.ABORTED); - readyQueue.remove(task.getDriverTaskId()); + if (readyQueue.remove(task.getDriverTaskId()) == null) { + readyQueue.decreaseReservedSize(task); + } break; case BLOCKED: task.setStatus(DriverTaskStatus.ABORTED); blockedTasks.remove(task); - readyQueue.decreaseReservedSize(); + readyQueue.decreaseReservedSize(task); break; case RUNNING: task.setStatus(DriverTaskStatus.ABORTED); - readyQueue.decreaseReservedSize(); + readyQueue.decreaseReservedSize(task); break; case FINISHED: break; @@ -489,6 +491,7 @@ public boolean readyToRunning(DriverTask task) { task.lock(); try { if (task.getStatus() != DriverTaskStatus.READY) { + readyQueue.decreaseReservedSize(task); return false; } @@ -545,7 +548,7 @@ public void runningToFinished(DriverTask task, ExecutionContext context) { } task.updateSchedulePriority(context); task.setStatus(DriverTaskStatus.FINISHED); - readyQueue.decreaseReservedSize(); + readyQueue.decreaseReservedSize(task); } finally { task.unlock(); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java index 1891fb6355677..aaf18685920a0 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java @@ -47,6 +47,7 @@ public synchronized E poll() throws InterruptedException { E output = pollFirst(); size--; reservedSize++; + markReserved(output); return output; } @@ -66,7 +67,7 @@ public synchronized void repush(E element) { throw new NullPointerException("pushed element is null"); } pushToQueue(element); - reservedSize--; + decreaseReservedSizeIfNecessary(element); size++; this.notifyAll(); } @@ -74,7 +75,32 @@ public synchronized void repush(E element) { /** * For task that is not in readyQueue when it's cleared, it won't be added into the queue again. */ - public synchronized void decreaseReservedSize() { - this.reservedSize--; + public synchronized boolean decreaseReservedSize(E element) { + if (element == null) { + throw new NullPointerException("pushed element is null"); + } + return decreaseReservedSizeIfNecessary(element); + } + + @Override + public synchronized void clear() { + super.clear(); + this.reservedSize = 0; + } + + protected void markReserved(E element) { + // Do nothing by default. + } + + protected boolean releaseReserved(E element) { + return true; + } + + private boolean decreaseReservedSizeIfNecessary(E element) { + if (!releaseReserved(element) || reservedSize <= 0) { + return false; + } + reservedSize--; + return true; } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/multilevelqueue/MultilevelPriorityQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/multilevelqueue/MultilevelPriorityQueue.java index c5f3cb3e6df1a..06c7275473b4e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/multilevelqueue/MultilevelPriorityQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/multilevelqueue/MultilevelPriorityQueue.java @@ -191,6 +191,16 @@ protected DriverTask get(DriverTask driverTask) { "MultilevelPriorityQueue does not support access element by get."); } + @Override + protected void markReserved(DriverTask task) { + task.markReservedInReadyQueue(); + } + + @Override + protected boolean releaseReserved(DriverTask task) { + return task.releaseReservedInReadyQueue(); + } + @Override protected void clearAllElements() { highestPriorityLevelQueue.clear(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/task/DriverTask.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/task/DriverTask.java index f320c6cad6c2f..85bada856caa2 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/task/DriverTask.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/task/DriverTask.java @@ -59,6 +59,7 @@ public class DriverTask implements IDIndexedAccessible { private final DriverTaskHandle driverTaskHandle; private long lastEnterReadyQueueTime; private long lastEnterBlockQueueTime; + private boolean reservedInReadyQueue; private long estimatedMemorySize; @@ -210,6 +211,18 @@ public void setLastEnterBlockQueueTime(long lastEnterBlockQueueTime) { this.lastEnterBlockQueueTime = lastEnterBlockQueueTime; } + public void markReservedInReadyQueue() { + reservedInReadyQueue = true; + } + + public boolean releaseReservedInReadyQueue() { + if (!reservedInReadyQueue) { + return false; + } + reservedInReadyQueue = false; + return true; + } + /** a comparator of ddl, the less the ddl is, the low order it has. */ public static class TimeoutComparator implements Comparator { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java index 0cc0e5edd71d5..0a5bb5ed3c04f 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java @@ -42,6 +42,7 @@ import org.mockito.Mockito; import java.io.IOException; +import java.lang.reflect.Field; import java.util.HashSet; import java.util.Map; import java.util.OptionalInt; @@ -197,6 +198,47 @@ public void testReadyToRunning() { clear(); } + @Test + public void testAbortReadyTaskAfterPollReleasesReadyQueueReservation() throws Exception { + IMPPDataExchangeManager mockMPPDataExchangeManager = + Mockito.mock(IMPPDataExchangeManager.class); + manager.setBlockManager(mockMPPDataExchangeManager); + ITaskScheduler defaultScheduler = manager.getScheduler(); + IDriver mockDriver = Mockito.mock(IDriver.class); + DriverTaskHandle driverTaskHandle = + new DriverTaskHandle( + 1, + (MultilevelPriorityQueue) manager.getReadyQueue(), + OptionalInt.of(Integer.MAX_VALUE)); + + QueryId queryId = new QueryId("test"); + FragmentInstanceId instanceId = + new FragmentInstanceId(new PlanFragmentId(queryId, 0), "inst-0"); + DriverTaskId driverTaskID = new DriverTaskId(instanceId, 0); + Mockito.when(mockDriver.getDriverTaskId()).thenReturn(driverTaskID); + DriverTask testTask = + new DriverTask(mockDriver, 100L, DriverTaskStatus.READY, driverTaskHandle, 0, false); + + try { + manager.registerTaskToQueryMap(queryId, testTask); + manager.submitTaskToReadyQueue(testTask); + + DriverTask polledTask = manager.getReadyQueue().poll(); + Assert.assertSame(testTask, polledTask); + Assert.assertEquals(1, getReadyQueueReservedSize()); + + manager.abortFragmentInstance(instanceId); + + Assert.assertEquals(DriverTaskStatus.ABORTED, testTask.getStatus()); + Assert.assertEquals(0, manager.getReadyQueue().size()); + Assert.assertEquals(0, getReadyQueueReservedSize()); + Assert.assertFalse(defaultScheduler.readyToRunning(polledTask)); + Assert.assertEquals(0, getReadyQueueReservedSize()); + } finally { + resetReadyQueueReservedSize(); + } + } + @Test public void testRunningToReady() { IMPPDataExchangeManager mockMPPDataExchangeManager = @@ -493,6 +535,33 @@ private void clear() { manager.getQueryMap().clear(); manager.getBlockedTasks().clear(); manager.getReadyQueue().clear(); + resetReadyQueueReservedSize(); manager.getTimeoutQueue().clear(); } + + private int getReadyQueueReservedSize() throws IllegalAccessException { + return getReadyQueueReservedSizeField().getInt(manager.getReadyQueue()); + } + + private void resetReadyQueueReservedSize() { + try { + getReadyQueueReservedSizeField().setInt(manager.getReadyQueue(), 0); + } catch (IllegalAccessException e) { + throw new AssertionError(e); + } + } + + private Field getReadyQueueReservedSizeField() { + Class queueClass = manager.getReadyQueue().getClass(); + while (queueClass != null) { + try { + Field reservedSizeField = queueClass.getDeclaredField("reservedSize"); + reservedSizeField.setAccessible(true); + return reservedSizeField; + } catch (NoSuchFieldException e) { + queueClass = queueClass.getSuperclass(); + } + } + throw new AssertionError("Cannot find reservedSize field in readyQueue hierarchy"); + } } From 9903696c0eddb7dd69518bf53cc5fd3f11206a9d Mon Sep 17 00:00:00 2001 From: JackieTien97 Date: Thu, 11 Jun 2026 18:18:16 +0800 Subject: [PATCH 2/2] Add ready queue reserved size metric --- .../execution/schedule/DriverScheduler.java | 4 +++ .../queue/IndexedBlockingReserveQueue.java | 4 +++ .../metric/DriverSchedulerMetricSet.java | 13 +++++++ .../schedule/DefaultDriverSchedulerTest.java | 36 +++---------------- 4 files changed, 25 insertions(+), 32 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java index 03110a8f1dd83..239ebb11329ae 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/DriverScheduler.java @@ -416,6 +416,10 @@ public long getReadyQueueTaskCount() { return readyQueue.size(); } + public long getReadyQueueReservedTaskCount() { + return readyQueue.getReservedSize(); + } + public long getBlockQueueTaskCount() { return blockedTasks.size(); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java index aaf18685920a0..b0f0edbd0e1aa 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/schedule/queue/IndexedBlockingReserveQueue.java @@ -82,6 +82,10 @@ public synchronized boolean decreaseReservedSize(E element) { return decreaseReservedSizeIfNecessary(element); } + public final synchronized int getReservedSize() { + return reservedSize; + } + @Override public synchronized void clear() { super.clear(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/metric/DriverSchedulerMetricSet.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/metric/DriverSchedulerMetricSet.java index 0b9e0c04f050d..94aaccfd9b7d6 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/metric/DriverSchedulerMetricSet.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/metric/DriverSchedulerMetricSet.java @@ -39,6 +39,7 @@ private DriverSchedulerMetricSet() { public static final String READY_QUEUED_TIME = "ready_queued_time"; public static final String BLOCK_QUEUED_TIME = "block_queued_time"; public static final String READY_QUEUE_TASK_COUNT = "ready_queue_task_count"; + public static final String READY_QUEUE_RESERVED_TASK_COUNT = "ready_queue_reserved_task_count"; public static final String BLOCK_QUEUE_TASK_COUNT = "block_queue_task_count"; private static final String TIMEOUT_QUEUE_SIZE = "timeout_queue_task_count"; private static final String QUERY_MAP_SIZE = "query_map_size"; @@ -67,6 +68,13 @@ public void bindTo(AbstractMetricService metricService) { DriverScheduler::getReadyQueueTaskCount, Tag.NAME.toString(), READY_QUEUE_TASK_COUNT); + metricService.createAutoGauge( + Metric.DRIVER_SCHEDULER.toString(), + MetricLevel.IMPORTANT, + DriverScheduler.getInstance(), + DriverScheduler::getReadyQueueReservedTaskCount, + Tag.NAME.toString(), + READY_QUEUE_RESERVED_TASK_COUNT); metricService.createAutoGauge( Metric.DRIVER_SCHEDULER.toString(), MetricLevel.IMPORTANT, @@ -109,6 +117,11 @@ public void unbindFrom(AbstractMetricService metricService) { Metric.DRIVER_SCHEDULER.toString(), Tag.NAME.toString(), READY_QUEUE_TASK_COUNT); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.DRIVER_SCHEDULER.toString(), + Tag.NAME.toString(), + READY_QUEUE_RESERVED_TASK_COUNT); metricService.remove( MetricType.AUTO_GAUGE, Metric.DRIVER_SCHEDULER.toString(), diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java index 0a5bb5ed3c04f..595d5febe4c6e 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/schedule/DefaultDriverSchedulerTest.java @@ -42,7 +42,6 @@ import org.mockito.Mockito; import java.io.IOException; -import java.lang.reflect.Field; import java.util.HashSet; import java.util.Map; import java.util.OptionalInt; @@ -225,17 +224,17 @@ public void testAbortReadyTaskAfterPollReleasesReadyQueueReservation() throws Ex DriverTask polledTask = manager.getReadyQueue().poll(); Assert.assertSame(testTask, polledTask); - Assert.assertEquals(1, getReadyQueueReservedSize()); + Assert.assertEquals(1, manager.getReadyQueueReservedTaskCount()); manager.abortFragmentInstance(instanceId); Assert.assertEquals(DriverTaskStatus.ABORTED, testTask.getStatus()); Assert.assertEquals(0, manager.getReadyQueue().size()); - Assert.assertEquals(0, getReadyQueueReservedSize()); + Assert.assertEquals(0, manager.getReadyQueueReservedTaskCount()); Assert.assertFalse(defaultScheduler.readyToRunning(polledTask)); - Assert.assertEquals(0, getReadyQueueReservedSize()); + Assert.assertEquals(0, manager.getReadyQueueReservedTaskCount()); } finally { - resetReadyQueueReservedSize(); + clear(); } } @@ -535,33 +534,6 @@ private void clear() { manager.getQueryMap().clear(); manager.getBlockedTasks().clear(); manager.getReadyQueue().clear(); - resetReadyQueueReservedSize(); manager.getTimeoutQueue().clear(); } - - private int getReadyQueueReservedSize() throws IllegalAccessException { - return getReadyQueueReservedSizeField().getInt(manager.getReadyQueue()); - } - - private void resetReadyQueueReservedSize() { - try { - getReadyQueueReservedSizeField().setInt(manager.getReadyQueue(), 0); - } catch (IllegalAccessException e) { - throw new AssertionError(e); - } - } - - private Field getReadyQueueReservedSizeField() { - Class queueClass = manager.getReadyQueue().getClass(); - while (queueClass != null) { - try { - Field reservedSizeField = queueClass.getDeclaredField("reservedSize"); - reservedSizeField.setAccessible(true); - return reservedSizeField; - } catch (NoSuchFieldException e) { - queueClass = queueClass.getSuperclass(); - } - } - throw new AssertionError("Cannot find reservedSize field in readyQueue hierarchy"); - } }