This is an automated email from the ASF dual-hosted git repository. jackietien pushed a commit to branch ChangeStatusCode in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit ce29a91e285d109263224e624c3f07e43114a5f0 Author: JackieTien97 <[email protected]> AuthorDate: Tue Jun 24 11:12:54 2025 +0800 Throw CANNOT_FETCH_FI_STATE(722) instead of 301/305 while DN restarting --- .../queryengine/execution/QueryStateMachine.java | 30 +++++++++---------- .../scheduler/FixedRateFragInsStateTracker.java | 34 +++++++++++----------- 2 files changed, 30 insertions(+), 34 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/QueryStateMachine.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/QueryStateMachine.java index b60d5e37e35..a6e3acfb44d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/QueryStateMachine.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/QueryStateMachine.java @@ -28,7 +28,6 @@ import org.apache.iotdb.rpc.RpcUtils; import com.google.common.util.concurrent.ListenableFuture; -import java.util.concurrent.Executor; import java.util.concurrent.ExecutorService; import static com.google.common.base.Preconditions.checkArgument; @@ -51,19 +50,13 @@ import static org.apache.iotdb.db.utils.ErrorHandlingUtils.getRootCause; public class QueryStateMachine { private final StateMachine<QueryState> queryState; - // The executor will be used in all the state machines belonged to this query. - private Executor stateMachineExecutor; private Throwable failureException; private TSStatus failureStatus; public QueryStateMachine(QueryId queryId, ExecutorService executor) { - this.stateMachineExecutor = executor; this.queryState = new StateMachine<>( - queryId.toString(), - this.stateMachineExecutor, - QUEUED, - QueryState.TERMINAL_INSTANCE_STATES); + queryId.toString(), executor, QUEUED, QueryState.TERMINAL_INSTANCE_STATES); } public void addStateChangeListener( @@ -112,9 +105,10 @@ public class QueryStateMachine { } public void transitionToCanceled(Throwable throwable, TSStatus failureStatus) { - this.failureException = throwable; - this.failureStatus = failureStatus; - transitionToDoneState(CANCELED); + if (transitionToDoneState(CANCELED)) { + this.failureException = throwable; + this.failureStatus = failureStatus; + } } public void transitionToAborted() { @@ -126,20 +120,22 @@ public class QueryStateMachine { } public void transitionToFailed(Throwable throwable) { - this.failureException = throwable; - transitionToDoneState(FAILED); + if (transitionToDoneState(FAILED)) { + this.failureException = throwable; + } } public void transitionToFailed(TSStatus failureStatus) { - this.failureStatus = failureStatus; - transitionToDoneState(FAILED); + if (transitionToDoneState(FAILED)) { + this.failureStatus = failureStatus; + } } - private void transitionToDoneState(QueryState doneState) { + private boolean transitionToDoneState(QueryState doneState) { requireNonNull(doneState, "doneState is null"); checkArgument(doneState.isDone(), "doneState %s is not a done state", doneState); - queryState.setIf(doneState, currentState -> !currentState.isDone()); + return queryState.setIf(doneState, currentState -> !currentState.isDone()); } public String getFailureMessage() { 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 f2afe5101ef..b848db6b6d9 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 @@ -180,8 +180,7 @@ public class FixedRateFragInsStateTracker extends AbstractFragInsStateTracker { instanceId, instanceInfo.getMessage()), TSStatusCode.CANNOT_FETCH_FI_STATE.getStatusCode(), true)); - } - if (instanceInfo.getState().isFailed()) { + } else if (instanceInfo.getState().isFailed()) { if (instanceInfo.getErrorCode().isPresent()) { stateMachine.transitionToFailed( new IoTDBException( @@ -196,22 +195,23 @@ public class FixedRateFragInsStateTracker extends AbstractFragInsStateTracker { } else { stateMachine.transitionToFailed(instanceInfo.getFailureInfoList().get(0).toException()); } - } - boolean queryFinished = false; - List<InstanceStateMetrics> rootInstanceStateMetricsList = - instanceStateMap.values().stream() - .filter(instanceStateMetrics -> instanceStateMetrics.isRootInstance) - .collect(Collectors.toList()); - if (!rootInstanceStateMetricsList.isEmpty()) { - queryFinished = - rootInstanceStateMetricsList.stream() - .allMatch( - instanceStateMetrics -> - instanceStateMetrics.lastState == FragmentInstanceState.FINISHED); - } + } else { + boolean queryFinished = false; + List<InstanceStateMetrics> rootInstanceStateMetricsList = + instanceStateMap.values().stream() + .filter(instanceStateMetrics -> instanceStateMetrics.isRootInstance) + .collect(Collectors.toList()); + if (!rootInstanceStateMetricsList.isEmpty()) { + queryFinished = + rootInstanceStateMetricsList.stream() + .allMatch( + instanceStateMetrics -> + instanceStateMetrics.lastState == FragmentInstanceState.FINISHED); + } - if (queryFinished) { - stateMachine.transitionToFinished(); + if (queryFinished) { + stateMachine.transitionToFinished(); + } } }
