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());