This is an automated email from the ASF dual-hosted git repository.
jackietien 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 3f00fdd6933 Cancel query which contains FI sent to shutdown datanode
by mistake
3f00fdd6933 is described below
commit 3f00fdd6933bec8fa25eaa8bba4b091c13709338
Author: Jackie Tien <[email protected]>
AuthorDate: Mon May 26 14:02:45 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 +++++++++++++++++-----
.../apache/iotdb/db/utils/ErrorHandlingUtils.java | 3 +-
5 files changed, 59 insertions(+), 16 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 7c212bd9548..d49fbd04bb1 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;
}
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/ErrorHandlingUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/ErrorHandlingUtils.java
index d8497e077d2..4cf091da012 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/ErrorHandlingUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/ErrorHandlingUtils.java
@@ -117,7 +117,8 @@ public class ErrorHandlingUtils {
|| status.getCode() ==
TSStatusCode.UDF_LOAD_CLASS_ERROR.getStatusCode()
|| status.getCode() ==
TSStatusCode.PLAN_FAILED_NETWORK_PARTITION.getStatusCode()
|| status.getCode() ==
TSStatusCode.SYNC_CONNECTION_ERROR.getStatusCode()
- || status.getCode() ==
TSStatusCode.NO_AVAILABLE_REPLICA.getStatusCode()) {
+ || status.getCode() ==
TSStatusCode.NO_AVAILABLE_REPLICA.getStatusCode()
+ || status.getCode() ==
TSStatusCode.CANNOT_FETCH_FI_STATE.getStatusCode()) {
LOGGER.info(message);
} else {
LOGGER.warn(message, e);