This is an automated email from the ASF dual-hosted git repository. shuwenwei pushed a commit to branch move-state-machine-to-calc-commons in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 1beb1668dd5c0c5534e4dee75653ad4555af9baf Author: shuwenwei <[email protected]> AuthorDate: Thu Jul 23 15:16:23 2026 +0800 Move state machine to calc commons --- .../org/apache/iotdb/calc/i18n/CalcMessages.java | 20 ++++++++ .../org/apache/iotdb/calc/i18n/CalcMessages.java | 20 ++++++++ .../iotdb/calc}/execution/FutureStateChange.java | 6 +-- .../apache/iotdb/calc}/execution/StateMachine.java | 55 ++++++++++------------ .../iotdb/db/i18n/DataNodeQueryMessages.java | 8 ---- .../iotdb/db/i18n/DataNodeQueryMessages.java | 11 ----- .../queryengine/execution/QueryStateMachine.java | 1 + .../fragment/FragmentInstanceStateMachine.java | 4 +- 8 files changed, 72 insertions(+), 53 deletions(-) diff --git a/iotdb-core/calc-commons/src/main/i18n/en/org/apache/iotdb/calc/i18n/CalcMessages.java b/iotdb-core/calc-commons/src/main/i18n/en/org/apache/iotdb/calc/i18n/CalcMessages.java index 9e22e34ceb9..d8ec4e0e6f9 100644 --- a/iotdb-core/calc-commons/src/main/i18n/en/org/apache/iotdb/calc/i18n/CalcMessages.java +++ b/iotdb-core/calc-commons/src/main/i18n/en/org/apache/iotdb/calc/i18n/CalcMessages.java @@ -437,6 +437,26 @@ public final class CalcMessages { "Error notifying state change listener for {}"; public static final String SERVER_IS_SHUTTING_DOWN = "Server is shutting down"; + public static final String EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4 = "executor is null"; + public static final String EXCEPTION_INITIALSTATE_IS_NULL_8992A39F = "initialState is null"; + public static final String EXCEPTION_TERMINALSTATES_IS_NULL_E0FC2A93 = "terminalStates is null"; + public static final String EXCEPTION_EXPECTEDSTATE_IS_NULL_5E8C2F32 = "expectedState is null"; + public static final String EXCEPTION_CURRENTSTATE_IS_NULL_AEDB20DB = "currentState is null"; + public static final String EXCEPTION_STATECHANGELISTENER_IS_NULL_635AE7D2 = + "stateChangeListener is null"; + public static final String EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30 = + "%s cannot transition from %s to %s"; + public static final String + EXCEPTION_CANNOT_FIRE_STATE_CHANGE_EVENT_WHILE_HOLDING_THE_LOCK_35243BC4 = + "Cannot fire state change event while holding the lock"; + public static final String EXCEPTION_CANNOT_NOTIFY_WHILE_HOLDING_THE_LOCK_15625D48 = + "Cannot notify while holding the lock"; + public static final String EXCEPTION_CANNOT_SET_STATE_WHILE_HOLDING_THE_LOCK_FA358188 = + "Cannot set state while holding the lock"; + public static final String + EXCEPTION_CANNOT_WAIT_FOR_STATE_CHANGE_WHILE_HOLDING_THE_LOCK_CBD9F784 = + "Cannot wait for state change while holding the lock"; + public static final String EXCEPTION_NEWSTATE_IS_NULL_D29A5454 = "newState is null"; public static final String EXCEPTION_PERCENTAGE_SHOULD_BE_IN_0_1_GOT_7A2C2F83 = "percentage should be in [0,1], got "; diff --git a/iotdb-core/calc-commons/src/main/i18n/zh/org/apache/iotdb/calc/i18n/CalcMessages.java b/iotdb-core/calc-commons/src/main/i18n/zh/org/apache/iotdb/calc/i18n/CalcMessages.java index 817d8376e60..6d8eb2bca0d 100644 --- a/iotdb-core/calc-commons/src/main/i18n/zh/org/apache/iotdb/calc/i18n/CalcMessages.java +++ b/iotdb-core/calc-commons/src/main/i18n/zh/org/apache/iotdb/calc/i18n/CalcMessages.java @@ -414,6 +414,26 @@ public final class CalcMessages { "通知 {} 的状态变更监听器时出错"; public static final String SERVER_IS_SHUTTING_DOWN = "服务器正在关闭"; + public static final String EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4 = "executor 不能为空"; + public static final String EXCEPTION_INITIALSTATE_IS_NULL_8992A39F = "initialState 不能为空"; + public static final String EXCEPTION_TERMINALSTATES_IS_NULL_E0FC2A93 = "terminalStates 不能为空"; + public static final String EXCEPTION_EXPECTEDSTATE_IS_NULL_5E8C2F32 = "expectedState 不能为空"; + public static final String EXCEPTION_CURRENTSTATE_IS_NULL_AEDB20DB = "currentState 不能为空"; + public static final String EXCEPTION_STATECHANGELISTENER_IS_NULL_635AE7D2 = + "stateChangeListener 不能为空"; + public static final String EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30 = + "%s 无法从 %s 转换到 %s"; + public static final String + EXCEPTION_CANNOT_FIRE_STATE_CHANGE_EVENT_WHILE_HOLDING_THE_LOCK_35243BC4 = + "持有锁时无法触发状态变更事件。"; + public static final String EXCEPTION_CANNOT_NOTIFY_WHILE_HOLDING_THE_LOCK_15625D48 = + "持有锁时无法通知。"; + public static final String EXCEPTION_CANNOT_SET_STATE_WHILE_HOLDING_THE_LOCK_FA358188 = + "持有锁时无法设置状态。"; + public static final String + EXCEPTION_CANNOT_WAIT_FOR_STATE_CHANGE_WHILE_HOLDING_THE_LOCK_CBD9F784 = + "持有锁时无法等待状态变更。"; + public static final String EXCEPTION_NEWSTATE_IS_NULL_D29A5454 = "newState 不能为空"; public static final String EXCEPTION_PERCENTAGE_SHOULD_BE_IN_0_1_GOT_7A2C2F83 = "百分比应在 [0,1] 范围内,实际为 "; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/FutureStateChange.java b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/FutureStateChange.java similarity index 92% rename from iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/FutureStateChange.java rename to iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/FutureStateChange.java index 1797fb9fe6b..aa1db91e177 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/FutureStateChange.java +++ b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/FutureStateChange.java @@ -17,9 +17,9 @@ * under the License. */ -package org.apache.iotdb.db.queryengine.execution; +package org.apache.iotdb.calc.execution; -import org.apache.iotdb.db.i18n.DataNodeQueryMessages; +import org.apache.iotdb.calc.i18n.CalcMessages; import com.google.common.collect.ImmutableSet; import com.google.common.util.concurrent.ListenableFuture; @@ -68,7 +68,7 @@ public class FutureStateChange<T> { } private void fireStateChange(T newState, Executor executor) { - requireNonNull(executor, DataNodeQueryMessages.EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4); + requireNonNull(executor, CalcMessages.EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4); Set<SettableFuture<T>> futures; synchronized (listeners) { futures = ImmutableSet.copyOf(listeners); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/StateMachine.java b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/StateMachine.java similarity index 85% rename from iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/StateMachine.java rename to iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/StateMachine.java index 3e3252e40af..43dc38b64b2 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/StateMachine.java +++ b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/StateMachine.java @@ -17,10 +17,9 @@ * under the License. */ -package org.apache.iotdb.db.queryengine.execution; +package org.apache.iotdb.calc.execution; import org.apache.iotdb.calc.i18n.CalcMessages; -import org.apache.iotdb.db.i18n.DataNodeQueryMessages; import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.ImmutableList; @@ -48,9 +47,6 @@ import static java.util.Objects.requireNonNull; @ThreadSafe public class StateMachine<T> { - private static final String LOCK_HELD_ERROR_MSG = "Cannot set state while holding the lock"; - private static final String STATE_IS_NULL = "newState is null"; - private static final Logger LOGGER = LoggerFactory.getLogger(StateMachine.class); private final String name; @@ -88,15 +84,12 @@ public class StateMachine<T> { * @param terminalStates the terminal states */ public StateMachine(String name, Executor executor, T initialState, Iterable<T> terminalStates) { - this.name = requireNonNull(name, DataNodeQueryMessages.EXCEPTION_NAME_IS_NULL_C8B35959); - this.executor = - requireNonNull(executor, DataNodeQueryMessages.EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4); - this.state = - requireNonNull(initialState, DataNodeQueryMessages.EXCEPTION_INITIALSTATE_IS_NULL_8992A39F); + this.name = requireNonNull(name, CalcMessages.EXCEPTION_NAME_IS_NULL_C8B35959); + this.executor = requireNonNull(executor, CalcMessages.EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4); + this.state = requireNonNull(initialState, CalcMessages.EXCEPTION_INITIALSTATE_IS_NULL_8992A39F); this.terminalStates = ImmutableSet.copyOf( - requireNonNull( - terminalStates, DataNodeQueryMessages.EXCEPTION_TERMINALSTATES_IS_NULL_E0FC2A93)); + requireNonNull(terminalStates, CalcMessages.EXCEPTION_TERMINALSTATES_IS_NULL_E0FC2A93)); } // State changes are atomic and state is volatile, so a direct read is safe here @@ -116,7 +109,7 @@ public class StateMachine<T> { T oldState = trySet(newState); checkState( oldState.equals(newState) || !isTerminalState(oldState), - DataNodeQueryMessages.EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30, + CalcMessages.EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30, name, state, newState); @@ -131,8 +124,10 @@ public class StateMachine<T> { * @return the state before the possible state change */ public T trySet(T newState) { - checkState(!Thread.holdsLock(lock), LOCK_HELD_ERROR_MSG); - requireNonNull(newState, STATE_IS_NULL); + checkState( + !Thread.holdsLock(lock), + CalcMessages.EXCEPTION_CANNOT_SET_STATE_WHILE_HOLDING_THE_LOCK_FA358188); + requireNonNull(newState, CalcMessages.EXCEPTION_NEWSTATE_IS_NULL_D29A5454); T oldState; FutureStateChange<T> oldFutureStateChange; @@ -166,8 +161,10 @@ public class StateMachine<T> { * @return true if the state is set */ public boolean setIf(T newState, Predicate<T> predicate) { - checkState(!Thread.holdsLock(lock), LOCK_HELD_ERROR_MSG); - requireNonNull(newState, STATE_IS_NULL); + checkState( + !Thread.holdsLock(lock), + CalcMessages.EXCEPTION_CANNOT_SET_STATE_WHILE_HOLDING_THE_LOCK_FA358188); + requireNonNull(newState, CalcMessages.EXCEPTION_NEWSTATE_IS_NULL_D29A5454); while (true) { // check if the current state passes the predicate @@ -197,9 +194,11 @@ public class StateMachine<T> { * @return true if the state is set */ public boolean compareAndSet(T expectedState, T newState) { - checkState(!Thread.holdsLock(lock), LOCK_HELD_ERROR_MSG); - requireNonNull(expectedState, DataNodeQueryMessages.EXCEPTION_EXPECTEDSTATE_IS_NULL_5E8C2F32); - requireNonNull(newState, STATE_IS_NULL); + checkState( + !Thread.holdsLock(lock), + CalcMessages.EXCEPTION_CANNOT_SET_STATE_WHILE_HOLDING_THE_LOCK_FA358188); + requireNonNull(expectedState, CalcMessages.EXCEPTION_EXPECTEDSTATE_IS_NULL_5E8C2F32); + requireNonNull(newState, CalcMessages.EXCEPTION_NEWSTATE_IS_NULL_D29A5454); FutureStateChange<T> oldFutureStateChange; ImmutableList<StateChangeListener<T>> curStateChangeListeners; @@ -215,7 +214,7 @@ public class StateMachine<T> { checkState( !isTerminalState(state), - DataNodeQueryMessages.EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30, + CalcMessages.EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30, name, state, newState); @@ -243,16 +242,15 @@ public class StateMachine<T> { List<StateChangeListener<T>> stateChangeListeners) { checkState( !Thread.holdsLock(lock), - DataNodeQueryMessages - .EXCEPTION_CANNOT_FIRE_STATE_CHANGE_EVENT_WHILE_HOLDING_THE_LOCK_35243BC4); - requireNonNull(newState, STATE_IS_NULL); + CalcMessages.EXCEPTION_CANNOT_FIRE_STATE_CHANGE_EVENT_WHILE_HOLDING_THE_LOCK_35243BC4); + requireNonNull(newState, CalcMessages.EXCEPTION_NEWSTATE_IS_NULL_D29A5454); // always fire listener callbacks from a different thread safeExecute( () -> { checkState( !Thread.holdsLock(lock), - DataNodeQueryMessages.EXCEPTION_CANNOT_NOTIFY_WHILE_HOLDING_THE_LOCK_15625D48); + CalcMessages.EXCEPTION_CANNOT_NOTIFY_WHILE_HOLDING_THE_LOCK_15625D48); try { futureStateChange.complete(newState); } catch (Throwable e) { @@ -280,9 +278,8 @@ public class StateMachine<T> { public ListenableFuture<T> getStateChange(T currentState) { checkState( !Thread.holdsLock(lock), - DataNodeQueryMessages - .EXCEPTION_CANNOT_WAIT_FOR_STATE_CHANGE_WHILE_HOLDING_THE_LOCK_CBD9F784); - requireNonNull(currentState, DataNodeQueryMessages.EXCEPTION_CURRENTSTATE_IS_NULL_AEDB20DB); + CalcMessages.EXCEPTION_CANNOT_WAIT_FOR_STATE_CHANGE_WHILE_HOLDING_THE_LOCK_CBD9F784); + requireNonNull(currentState, CalcMessages.EXCEPTION_CURRENTSTATE_IS_NULL_AEDB20DB); synchronized (lock) { // return a completed future if the state has already changed, or we are in a terminal state @@ -303,7 +300,7 @@ public class StateMachine<T> { */ public void addStateChangeListener(StateChangeListener<T> stateChangeListener) { requireNonNull( - stateChangeListener, DataNodeQueryMessages.EXCEPTION_STATECHANGELISTENER_IS_NULL_635AE7D2); + stateChangeListener, CalcMessages.EXCEPTION_STATECHANGELISTENER_IS_NULL_635AE7D2); T currentState; synchronized (lock) { diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java index b540ebb2540..92242b62aa7 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java @@ -3292,15 +3292,7 @@ public final class DataNodeQueryMessages { public static final String EXCEPTION_INVALID_ARG_ARG_2946DBE5 = "Invalid %s %s"; public static final String EXCEPTION_ID_IS_EMPTY_28C94FC0 = "id is empty"; public static final String EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4 = "executor is null"; - public static final String EXCEPTION_INITIALSTATE_IS_NULL_8992A39F = "initialState is null"; - public static final String EXCEPTION_TERMINALSTATES_IS_NULL_E0FC2A93 = "terminalStates is null"; - public static final String EXCEPTION_EXPECTEDSTATE_IS_NULL_5E8C2F32 = "expectedState is null"; public static final String EXCEPTION_CURRENTSTATE_IS_NULL_AEDB20DB = "currentState is null"; - public static final String EXCEPTION_STATECHANGELISTENER_IS_NULL_635AE7D2 = "stateChangeListener is null"; - public static final String EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30 = "%s cannot transition from %s to %s"; - public static final String EXCEPTION_CANNOT_FIRE_STATE_CHANGE_EVENT_WHILE_HOLDING_THE_LOCK_35243BC4 = "Cannot fire state change event while holding the lock"; - public static final String EXCEPTION_CANNOT_NOTIFY_WHILE_HOLDING_THE_LOCK_15625D48 = "Cannot notify while holding the lock"; - public static final String EXCEPTION_CANNOT_WAIT_FOR_STATE_CHANGE_WHILE_HOLDING_THE_LOCK_CBD9F784 = "Cannot wait for state change while holding the lock"; public static final String EXCEPTION_DONESTATE_IS_NULL_D88F77E5 = "doneState is null"; public static final String EXCEPTION_DONESTATE_ARG_IS_NOT_A_DONE_STATE_8724C618 = "doneState %s is not a done state"; public static final String EXCEPTION_DATANODEID_SHOULD_BE_INIT_FIRST_13B19A85 = "DataNodeId should be init first!"; diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java index 0e73bd1dd96..9ce80200a7b 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java @@ -3960,18 +3960,7 @@ public final class DataNodeQueryMessages { public static final String EXCEPTION_ID_IS_EMPTY_28C94FC0 = "id 为空"; public static final String EXCEPTION_EXECUTOR_IS_NULL_7FBE03A4 = "executor 不能为空"; - public static final String EXCEPTION_INITIALSTATE_IS_NULL_8992A39F = "initialState 不能为空"; - public static final String EXCEPTION_TERMINALSTATES_IS_NULL_E0FC2A93 = "terminalStates 不能为空"; - public static final String EXCEPTION_EXPECTEDSTATE_IS_NULL_5E8C2F32 = "expectedState 不能为空"; public static final String EXCEPTION_CURRENTSTATE_IS_NULL_AEDB20DB = "currentState 不能为空"; - public static final String EXCEPTION_STATECHANGELISTENER_IS_NULL_635AE7D2 = "stateChangeListener 不能为空"; - public static final String EXCEPTION_ARG_CANNOT_TRANSITION_FROM_ARG_TO_ARG_8C680D30 = "%s 无法从 %s 转换到 %s"; - public static final String EXCEPTION_CANNOT_FIRE_STATE_CHANGE_EVENT_WHILE_HOLDING_THE_LOCK_35243BC4 = - "持有锁时无法触发状态变更事件。"; - public static final String EXCEPTION_CANNOT_NOTIFY_WHILE_HOLDING_THE_LOCK_15625D48 = - "持有锁时无法通知。"; - public static final String EXCEPTION_CANNOT_WAIT_FOR_STATE_CHANGE_WHILE_HOLDING_THE_LOCK_CBD9F784 = - "持有锁时无法等待状态变更。"; public static final String EXCEPTION_DONESTATE_IS_NULL_D88F77E5 = "doneState 不能为空"; public static final String EXCEPTION_DONESTATE_ARG_IS_NOT_A_DONE_STATE_8724C618 = "doneState %s 不是完成状态"; 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 e83b21979ea..1bc3e2b0f32 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 @@ -19,6 +19,7 @@ package org.apache.iotdb.db.queryengine.execution; +import org.apache.iotdb.calc.execution.StateMachine; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.exception.IoTDBException; import org.apache.iotdb.commons.exception.IoTDBRuntimeException; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceStateMachine.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceStateMachine.java index 92fd00db352..771d96693d4 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceStateMachine.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceStateMachine.java @@ -19,10 +19,10 @@ package org.apache.iotdb.db.queryengine.execution.fragment; +import org.apache.iotdb.calc.execution.StateMachine; +import org.apache.iotdb.calc.execution.StateMachine.StateChangeListener; import org.apache.iotdb.db.i18n.DataNodeQueryMessages; import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; -import org.apache.iotdb.db.queryengine.execution.StateMachine; -import org.apache.iotdb.db.queryengine.execution.StateMachine.StateChangeListener; import org.apache.iotdb.db.utils.SetThreadName; import com.google.common.collect.ImmutableList;
