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

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

commit d37a689c524778d37f2a40e4f12c2f251a373893
Author: JackieTien97 <[email protected]>
AuthorDate: Mon May 26 11:55:34 2025 +0800

    Cancel query which contains FI sent to shutdown datanode by mistake
---
 .../java/org/apache/iotdb/rpc/TSStatusCode.java    |  1 +
 .../execution/fragment/FragmentInstanceInfo.java   |  4 ++
 .../execution/fragment/FragmentInstanceState.java  |  2 +-
 .../scheduler/FixedRateFragInsStateTracker.java    | 65 +++++++++++++++++-----
 4 files changed, 57 insertions(+), 15 deletions(-)

diff --git 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
index 3c61ea03326..df4d141e203 100644
--- 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
+++ 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
@@ -140,6 +140,7 @@ public enum TSStatusCode {
   QUERY_EXECUTION_MEMORY_NOT_ENOUGH(719),
   QUERY_TIMEOUT(720),
   PLAN_FAILED_NETWORK_PARTITION(721),
+  CANNOT_FETCH_FI_STATE(722),
 
   // Arithmetic
   NUMERIC_VALUE_OUT_OF_RANGE(750),
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceInfo.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceInfo.java
index 4717a23f279..a544aebe6df 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceInfo.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceInfo.java
@@ -77,6 +77,10 @@ public class FragmentInstanceInfo implements DataSet {
     return message;
   }
 
+  public void setMessage(String message) {
+    this.message = message;
+  }
+
   public Optional<TSStatus> getErrorCode() {
     return Optional.ofNullable(errorCode);
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceState.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceState.java
index 092b1be3816..2bcb12544d9 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceState.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceState.java
@@ -47,7 +47,7 @@ public enum FragmentInstanceState {
   /** Instance execution failed. */
   FAILED(true, true),
   /** Instance is not found. */
-  NO_SUCH_INSTANCE(false, true);
+  NO_SUCH_INSTANCE(true, true);
 
   public static final Set<FragmentInstanceState> TERMINAL_INSTANCE_STATES =
       Stream.of(FragmentInstanceState.values())
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FixedRateFragInsStateTracker.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FixedRateFragInsStateTracker.java
index e18c628e421..f2afe5101ef 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FixedRateFragInsStateTracker.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FixedRateFragInsStateTracker.java
@@ -31,6 +31,7 @@ import 
org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceInfo;
 import 
org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceState;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.FragmentInstance;
 import org.apache.iotdb.db.utils.SetThreadName;
+import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.apache.thrift.TException;
 import org.slf4j.Logger;
@@ -45,13 +46,15 @@ import java.util.concurrent.ScheduledFuture;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 
+import static 
org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceState.NO_SUCH_INSTANCE;
+
 public class FixedRateFragInsStateTracker extends AbstractFragInsStateTracker {
 
   private static final Logger logger = 
LoggerFactory.getLogger(FixedRateFragInsStateTracker.class);
 
   private static final long SAME_STATE_PRINT_RATE_IN_MS = 10L * 60 * 1000;
 
-  // TODO: (xingtanzjr) consider how much Interval is OK for state tracker
+  // consider how much Interval is OK for state tracker
   private static final long STATE_FETCH_INTERVAL_IN_MS = 500;
   private ScheduledFuture<?> trackTask;
   private final Map<FragmentInstanceId, InstanceStateMetrics> instanceStateMap;
@@ -112,8 +115,8 @@ public class FixedRateFragInsStateTracker extends 
AbstractFragInsStateTracker {
     aborted = true;
     if (trackTask != null) {
       boolean cancelResult = trackTask.cancel(true);
-      // TODO: (xingtanzjr) a strange case here is that sometimes
-      // the cancelResult is false but the trackTask is definitely cancelled
+      // a strange case here is that sometimes the cancelResult is false but 
the trackTask is
+      // definitely cancelled
       if (!cancelResult) {
         logger.debug("cancel state tracking task failed. {}", 
trackTask.isCancelled());
       }
@@ -144,8 +147,24 @@ public class FixedRateFragInsStateTracker extends 
AbstractFragInsStateTracker {
             updateQueryState(instance.getId(), instanceInfo);
           }
         } catch (ClientManagerException | TException e) {
-          // TODO: do nothing ?
-          logger.warn("error happened while fetching query state", e);
+          // network exception, should retry
+          InstanceStateMetrics metrics =
+              instanceStateMap.computeIfAbsent(
+                  instance.getId(), k -> new 
InstanceStateMetrics(instance.isRoot()));
+          if (metrics.reachMaxRetryCount()) {
+            // if reach max retry count, we think that the DN is down, and FI 
in that node won't
+            // exist
+            FragmentInstanceInfo instanceInfo = new 
FragmentInstanceInfo(NO_SUCH_INSTANCE);
+            instanceInfo.setMessage(
+                String.format(
+                    "Failed to fetch state, has retried %s times",
+                    InstanceStateMetrics.MAX_STATE_FETCH_RETRY_COUNT));
+            updateQueryState(instance.getId(), instanceInfo);
+          } else {
+            // if not reaching max retry count, add retry count, and wait for 
next fetching schedule
+            metrics.addRetryCount();
+            logger.warn("error happened while fetching query state", e);
+          }
         }
       }
     }
@@ -153,25 +172,27 @@ public class FixedRateFragInsStateTracker extends 
AbstractFragInsStateTracker {
 
   private void updateQueryState(FragmentInstanceId instanceId, 
FragmentInstanceInfo instanceInfo) {
     // no such instance may be caused by DN restarting
-    if (instanceInfo.getState() == FragmentInstanceState.NO_SUCH_INSTANCE) {
+    if (instanceInfo.getState() == NO_SUCH_INSTANCE) {
       stateMachine.transitionToFailed(
-          new RuntimeException(
+          new IoTDBException(
               String.format(
                   "FragmentInstance[%s] is failed. %s, may be caused by DN 
restarting.",
-                  instanceId, instanceInfo.getMessage())));
+                  instanceId, instanceInfo.getMessage()),
+              TSStatusCode.CANNOT_FETCH_FI_STATE.getStatusCode(),
+              true));
     }
     if (instanceInfo.getState().isFailed()) {
-      if (instanceInfo.getFailureInfoList() == null
+      if (instanceInfo.getErrorCode().isPresent()) {
+        stateMachine.transitionToFailed(
+            new IoTDBException(
+                instanceInfo.getErrorCode().get().getMessage(),
+                instanceInfo.getErrorCode().get().getCode()));
+      } else if (instanceInfo.getFailureInfoList() == null
           || instanceInfo.getFailureInfoList().isEmpty()) {
         stateMachine.transitionToFailed(
             new RuntimeException(
                 String.format(
                     "FragmentInstance[%s] is failed. %s", instanceId, 
instanceInfo.getMessage())));
-      } else if (instanceInfo.getErrorCode().isPresent()) {
-        stateMachine.transitionToFailed(
-            new IoTDBException(
-                instanceInfo.getErrorCode().get().getMessage(),
-                instanceInfo.getErrorCode().get().getCode()));
       } else {
         
stateMachine.transitionToFailed(instanceInfo.getFailureInfoList().get(0).toException());
       }
@@ -203,23 +224,39 @@ public class FixedRateFragInsStateTracker extends 
AbstractFragInsStateTracker {
   }
 
   private static class InstanceStateMetrics {
+    private static final long MAX_STATE_FETCH_RETRY_COUNT = 5;
     private final boolean isRootInstance;
     private FragmentInstanceState lastState;
     private long durationToLastPrintInMS;
+    // we only record the continuous retry count
+    private int retryCount;
 
     private InstanceStateMetrics(boolean isRootInstance) {
       this.isRootInstance = isRootInstance;
       this.lastState = null;
       this.durationToLastPrintInMS = 0L;
+      this.retryCount = 0;
     }
 
     private void reset(FragmentInstanceState newState) {
       this.lastState = newState;
       this.durationToLastPrintInMS = 0L;
+      // each successful fetch, we need to reset the retry count
+      this.retryCount = 0;
+    }
+
+    private void addRetryCount() {
+      this.retryCount++;
+    }
+
+    private boolean reachMaxRetryCount() {
+      return retryCount >= MAX_STATE_FETCH_RETRY_COUNT;
     }
 
     private void addDuration(long duration) {
       durationToLastPrintInMS += duration;
+      // each successful fetch, we need to reset the retry count
+      this.retryCount = 0;
     }
   }
 }

Reply via email to