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

Reply via email to