This is an automated email from the ASF dual-hosted git repository. lancelly pushed a commit to branch priorityForShowQueries in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit c5acdf13cc63a4be580a58fbc03fbbbaa3b7e0ad Author: lancelly <[email protected]> AuthorDate: Fri Oct 13 09:57:06 2023 +0800 tmp --- .../queue/multilevelqueue/MultilevelPriorityQueue.java | 15 +++++++++++++++ .../queryengine/execution/schedule/task/DriverTask.java | 6 ++++++ 2 files changed, 21 insertions(+) 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 38e7c51e88b..6094410d74b 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 @@ -43,6 +43,10 @@ public class MultilevelPriorityQueue extends IndexedBlockingReserveQueue<DriverT private final PriorityQueue<DriverTask>[] levelWaitingSplits; + /** + This queue is independent of the other priority queues and has the highest priority. It is used to assign the highest execution priority to tasks like "ShowQuery," without considering cumulative execution time. */ + private final PriorityQueue<DriverTask> highestPriorityLevelQueue; + /** * Total amount of time each LEVEL has occupied, which decides which level we will take task from. */ @@ -65,6 +69,7 @@ public class MultilevelPriorityQueue extends IndexedBlockingReserveQueue<DriverT this.levelScheduledTime = new AtomicLong[LEVEL_THRESHOLD_SECONDS.length]; this.levelMinScheduledTime = new AtomicLong[LEVEL_THRESHOLD_SECONDS.length]; this.levelWaitingSplits = new PriorityQueue[LEVEL_THRESHOLD_SECONDS.length]; + this.highestPriorityLevelQueue = new PriorityQueue<>(); for (int level = 0; level < LEVEL_THRESHOLD_SECONDS.length; level++) { levelScheduledTime[level] = new AtomicLong(); levelMinScheduledTime[level] = new AtomicLong(-1); @@ -86,6 +91,11 @@ public class MultilevelPriorityQueue extends IndexedBlockingReserveQueue<DriverT @Override public void pushToQueue(DriverTask task) { checkArgument(task != null, "DriverTask to be pushed is null"); + // Push tasks with the highest priority(Currently, only ShowQuery related tasks) into highestPriorityLevelQueue directly. + if(task.isHighestPriority()){ + highestPriorityLevelQueue.offer(task); + return; + } int level = task.getPriority().getLevel(); if (levelWaitingSplits[level].isEmpty()) { @@ -102,6 +112,11 @@ public class MultilevelPriorityQueue extends IndexedBlockingReserveQueue<DriverT } protected DriverTask pollFirst() { + // Always choose tasks in the highestPriorityLevelQueue first. + if(!highestPriorityLevelQueue.isEmpty()){ + return highestPriorityLevelQueue.poll(); + } + DriverTask result; while (true) { result = chooseLevelAndTask(); 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 a3b616f42c4..f3ed7d29b69 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 @@ -50,6 +50,8 @@ public class DriverTask implements IDIndexedAccessible { private final long ddl; private final Lock lock; + private final boolean isHighestPriority = true; + private String abortCause; private final AtomicReference<Priority> priority; @@ -131,6 +133,10 @@ public class DriverTask implements IDIndexedAccessible { return ddl; } + public boolean isHighestPriority() { + return isHighestPriority; + } + @Override public int hashCode() { return driver.getDriverTaskId().hashCode();
