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

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


The following commit(s) were added to refs/heads/master by this push:
     new a7aa8ca238 [IOTDB-5450] Add only one task of a query to timeoutQueue 
(#8960)
a7aa8ca238 is described below

commit a7aa8ca2383a398194e2404d1e991be93336d92f
Author: Xiangwei Wei <[email protected]>
AuthorDate: Mon Feb 6 13:34:03 2023 +0800

    [IOTDB-5450] Add only one task of a query to timeoutQueue (#8960)
---
 .../db/mpp/execution/schedule/DriverScheduler.java | 23 +++++++----------
 .../mpp/execution/schedule/IDriverScheduler.java   | 11 --------
 .../execution/schedule/DriverSchedulerTest.java    | 29 ++++++++++++----------
 3 files changed, 25 insertions(+), 38 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
index ddcf072a1b..fb0fb7bf0c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
@@ -32,9 +32,7 @@ import 
org.apache.iotdb.db.mpp.execution.schedule.queue.IndexedBlockingQueue;
 import org.apache.iotdb.db.mpp.execution.schedule.queue.L1PriorityQueue;
 import 
org.apache.iotdb.db.mpp.execution.schedule.queue.multilevelqueue.DriverTaskHandle;
 import 
org.apache.iotdb.db.mpp.execution.schedule.queue.multilevelqueue.MultilevelPriorityQueue;
-import 
org.apache.iotdb.db.mpp.execution.schedule.queue.multilevelqueue.Priority;
 import org.apache.iotdb.db.mpp.execution.schedule.task.DriverTask;
-import org.apache.iotdb.db.mpp.execution.schedule.task.DriverTaskId;
 import org.apache.iotdb.db.mpp.execution.schedule.task.DriverTaskStatus;
 import org.apache.iotdb.db.mpp.metric.QueryMetricsManager;
 import org.apache.iotdb.db.utils.SetThreadName;
@@ -184,21 +182,28 @@ public class DriverScheduler implements IDriverScheduler, 
IService {
                         DriverTaskStatus.READY,
                         driverTaskHandle))
             .collect(Collectors.toList());
+    // If query has not been registered by other fragment instances,
+    // add the first task as timeout checking task to timeoutQueue.
     for (DriverTask driverTask : tasks) {
       queryMap
-          .computeIfAbsent(queryId, v -> new ConcurrentHashMap<>())
+          .computeIfAbsent(
+              queryId,
+              v -> {
+                timeoutQueue.push(tasks.get(0));
+                return new ConcurrentHashMap<>();
+              })
           .computeIfAbsent(
               driverTask.getDriverTaskId().getFragmentInstanceId(),
               v -> Collections.synchronizedSet(new HashSet<>()))
           .add(driverTask);
     }
+
     for (DriverTask task : tasks) {
       task.lock();
       try {
         if (task.getStatus() != DriverTaskStatus.READY) {
           continue;
         }
-        timeoutQueue.push(task);
         readyQueue.push(task);
         task.setLastEnterReadyQueueTime(System.nanoTime());
       } finally {
@@ -246,16 +251,6 @@ public class DriverScheduler implements IDriverScheduler, 
IService {
     }
   }
 
-  @Override
-  public Priority getSchedulePriority(DriverTaskId driverTaskID) {
-    DriverTask task = timeoutQueue.get(driverTaskID);
-    if (task == null) {
-      throw new IllegalStateException(
-          "the fragmentInstance " + driverTaskID.getFullId() + " has been 
cleared");
-    }
-    return task.getPriority();
-  }
-
   private void clearDriverTask(DriverTask task) {
     try (SetThreadName driverTaskName =
         new SetThreadName(task.getDriver().getDriverTaskId().getFullId())) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/IDriverScheduler.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/IDriverScheduler.java
index 5729a2718c..ba2259a200 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/IDriverScheduler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/IDriverScheduler.java
@@ -21,8 +21,6 @@ package org.apache.iotdb.db.mpp.execution.schedule;
 import org.apache.iotdb.db.mpp.common.FragmentInstanceId;
 import org.apache.iotdb.db.mpp.common.QueryId;
 import org.apache.iotdb.db.mpp.execution.driver.IDriver;
-import 
org.apache.iotdb.db.mpp.execution.schedule.queue.multilevelqueue.Priority;
-import org.apache.iotdb.db.mpp.execution.schedule.task.DriverTaskId;
 
 import java.util.List;
 
@@ -52,13 +50,4 @@ public interface IDriverScheduler {
    * @param instanceId the id of the fragment instance to be aborted.
    */
   void abortFragmentInstance(FragmentInstanceId instanceId);
-
-  /**
-   * Return the schedule priority of a Driver task.
-   *
-   * @param driverTaskID the fragment instance id.
-   * @return the schedule priority.
-   * @throws IllegalStateException if the instance has already been cleared.
-   */
-  Priority getSchedulePriority(DriverTaskId driverTaskID);
 }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/schedule/DriverSchedulerTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/schedule/DriverSchedulerTest.java
index 757f39358e..7e3612e499 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/schedule/DriverSchedulerTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/schedule/DriverSchedulerTest.java
@@ -73,14 +73,16 @@ public class DriverSchedulerTest {
     Assert.assertEquals(1, manager.getQueryMap().size());
     Assert.assertTrue(manager.getQueryMap().containsKey(queryId));
     Assert.assertEquals(2, manager.getQueryMap().get(queryId).size());
-    Assert.assertEquals(2, manager.getTimeoutQueue().size());
+    Assert.assertEquals(1, manager.getTimeoutQueue().size());
     Assert.assertEquals(2, manager.getReadyQueue().size());
-    DriverTask task1 = manager.getTimeoutQueue().get(driverTaskId1);
-    Assert.assertNotNull(task1);
-    DriverTask task2 = manager.getTimeoutQueue().get(driverTaskId2);
-    Assert.assertNotNull(task2);
-    
Assert.assertTrue(manager.getQueryMap().get(queryId).get(instanceId1).contains(task1));
-    
Assert.assertTrue(manager.getQueryMap().get(queryId).get(instanceId2).contains(task2));
+    Assert.assertNotNull(manager.getTimeoutQueue().get(driverTaskId1));
+    Assert.assertNull(manager.getTimeoutQueue().get(driverTaskId2));
+    DriverTask task1 =
+        (DriverTask) 
manager.getQueryMap().get(queryId).get(instanceId1).toArray()[0];
+    DriverTask task2 =
+        (DriverTask) 
manager.getQueryMap().get(queryId).get(instanceId2).toArray()[0];
+    Assert.assertEquals(task1.getDriverTaskId(), driverTaskId1);
+    Assert.assertEquals(task2.getDriverTaskId(), driverTaskId2);
     Assert.assertEquals(DriverTaskStatus.READY, task1.getStatus());
     Assert.assertEquals(DriverTaskStatus.READY, task2.getStatus());
 
@@ -94,11 +96,12 @@ public class DriverSchedulerTest {
     Assert.assertEquals(1, manager.getQueryMap().size());
     Assert.assertTrue(manager.getQueryMap().containsKey(queryId));
     Assert.assertEquals(3, manager.getQueryMap().get(queryId).size());
-    Assert.assertEquals(3, manager.getTimeoutQueue().size());
+    Assert.assertEquals(1, manager.getTimeoutQueue().size());
     Assert.assertEquals(3, manager.getReadyQueue().size());
-    DriverTask task3 = manager.getTimeoutQueue().get(driverTaskId3);
-    Assert.assertNotNull(task3);
-    
Assert.assertTrue(manager.getQueryMap().get(queryId).get(instanceId3).contains(task3));
+    Assert.assertNull(manager.getTimeoutQueue().get(driverTaskId3));
+    DriverTask task3 =
+        (DriverTask) 
manager.getQueryMap().get(queryId).get(instanceId3).toArray()[0];
+    Assert.assertEquals(task3.getDriverTaskId(), driverTaskId3);
     Assert.assertEquals(DriverTaskStatus.READY, task3.getStatus());
 
     // Submit another task of the different query
@@ -113,7 +116,7 @@ public class DriverSchedulerTest {
     Assert.assertEquals(2, manager.getQueryMap().size());
     Assert.assertTrue(manager.getQueryMap().containsKey(queryId2));
     Assert.assertEquals(1, manager.getQueryMap().get(queryId2).size());
-    Assert.assertEquals(4, manager.getTimeoutQueue().size());
+    Assert.assertEquals(2, manager.getTimeoutQueue().size());
     Assert.assertEquals(4, manager.getReadyQueue().size());
     DriverTask task4 = manager.getTimeoutQueue().get(driverTaskId4);
     Assert.assertNotNull(task4);
@@ -129,7 +132,7 @@ public class DriverSchedulerTest {
     Assert.assertTrue(manager.getBlockedTasks().isEmpty());
     Assert.assertEquals(2, manager.getQueryMap().size());
     Assert.assertTrue(manager.getQueryMap().containsKey(queryId));
-    Assert.assertEquals(3, manager.getTimeoutQueue().size());
+    Assert.assertEquals(1, manager.getTimeoutQueue().size());
     Assert.assertEquals(3, manager.getReadyQueue().size());
     Assert.assertEquals(DriverTaskStatus.ABORTED, task1.getStatus());
     Assert.assertEquals(DriverTaskStatus.READY, task2.getStatus());

Reply via email to