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