This is an automated email from the ASF dual-hosted git repository. xingtanzjr pushed a commit to branch xingtanzjr/fix_query_map in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 6b5ff8a3fc9b0c95865046cf3886b97f193869fe Author: Jinrui.Zhang <[email protected]> AuthorDate: Tue Jun 14 16:48:39 2022 +0800 add QueryExecution Remove method --- server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java | 4 ++++ .../org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java | 2 ++ .../apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java | 1 + 3 files changed, 7 insertions(+) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java index 872d8c3cfa..b8ff45b639 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java @@ -137,6 +137,10 @@ public class Coordinator { return queryExecutionMap.get(queryId); } + public void removeQueryExecution(Long queryId) { + queryExecutionMap.remove(queryId); + } + // TODO: (xingtanzjr) need to redo once we have a concrete policy for the threadPool management private ExecutorService getQueryExecutor() { return IoTDBThreadPoolFactory.newFixedThreadPool( diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java index 3ed55d9995..3568009de5 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java @@ -126,6 +126,8 @@ public class ClusterSchemaFetcher implements ISchemaFetcher { result.mergeSchemaTree(fetchedSchemaTree); } } + coordinator.getQueryExecution(queryId).stopAndCleanup(); + coordinator.removeQueryExecution(queryId); return result; } } diff --git a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java index 6e5ac4349e..0726fe4cfd 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java @@ -1304,6 +1304,7 @@ public class DataNodeTSIServiceImpl implements TSIEventHandler { try (SetThreadName threadName = new SetThreadName(queryExecution.getQueryId())) { LOGGER.info("stop and clean up"); queryExecution.stopAndCleanup(); + COORDINATOR.removeQueryExecution(queryId); } } }
