This is an automated email from the ASF dual-hosted git repository.
tanxinyu pushed a commit to branch rc/1.3.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rc/1.3.3 by this push:
new ec640661f6a [to rc/1.3.3] Split CnToDnRequestType to sync and async &
Add check for adding new request type (see #13660) (#13784)
ec640661f6a is described below
commit ec640661f6a383347091f733854fc1a3ebb08d5c
Author: Li Yu Heng <[email protected]>
AuthorDate: Mon Oct 28 11:49:15 2024 +0800
[to rc/1.3.3] Split CnToDnRequestType to sync and async & Add check for
adding new request type (see #13660) (#13784)
---
.../CnToDnAsyncRequestType.java} | 24 +---
.../CnToDnInternalServiceAsyncRequestManager.java | 144 ++++++++++++---------
.../handlers/DataNodeAsyncRequestContext.java | 12 +-
.../rpc/CheckTimeSeriesExistenceRPCHandler.java | 4 +-
.../rpc/CountPathsUsingTemplateRPCHandler.java | 4 +-
.../rpc/DataNodeAsyncRequestRPCHandler.java | 10 +-
.../handlers/rpc/DataNodeTSStatusRPCHandler.java | 4 +-
.../rpc/FetchSchemaBlackListRPCHandler.java | 4 +-
.../handlers/rpc/PipeHeartbeatRPCHandler.java | 4 +-
.../async/handlers/rpc/PipePushMetaRPCHandler.java | 4 +-
.../async/handlers/rpc/SchemaUpdateRPCHandler.java | 4 +-
.../rpc/SubmitTestConnectionTaskRPCHandler.java | 4 +-
.../handlers/rpc/TransferLeaderRPCHandler.java | 4 +-
.../CheckSchemaRegionUsingTemplateRPCHandler.java | 4 +-
.../ConsumerGroupPushMetaRPCHandler.java | 4 +-
.../rpc/subscription/TopicPushMetaRPCHandler.java | 4 +-
.../client/sync/CnToDnSyncRequestType.java | 49 +++++++
.../client/sync/SyncDataNodeClientPool.java | 136 ++++++++++++-------
.../confignode/conf/ConfigNodeStartupCheck.java | 9 ++
.../iotdb/confignode/manager/ClusterManager.java | 9 +-
.../confignode/manager/ClusterQuotaManager.java | 6 +-
.../iotdb/confignode/manager/TriggerManager.java | 4 +-
.../iotdb/confignode/manager/UDFManager.java | 6 +-
.../manager/load/balancer/RouteBalancer.java | 8 +-
.../iotdb/confignode/manager/node/NodeManager.java | 29 +++--
.../manager/partition/PartitionManager.java | 8 +-
.../runtime/heartbeat/PipeHeartbeatScheduler.java | 4 +-
.../manager/schema/ClusterSchemaManager.java | 4 +-
.../procedure/env/ConfigNodeProcedureEnv.java | 57 ++++----
.../procedure/env/RegionMaintainHandler.java | 22 ++--
.../impl/schema/AlterLogicalViewProcedure.java | 8 +-
.../impl/schema/DataNodeRegionTaskExecutor.java | 8 +-
.../impl/schema/DeactivateTemplateProcedure.java | 16 +--
.../impl/schema/DeleteDatabaseProcedure.java | 4 +-
.../impl/schema/DeleteLogicalViewProcedure.java | 12 +-
.../impl/schema/DeleteTimeSeriesProcedure.java | 16 +--
.../procedure/impl/schema/SchemaUtils.java | 6 +-
.../procedure/impl/schema/SetTTLProcedure.java | 4 +-
.../impl/schema/SetTemplateProcedure.java | 12 +-
.../impl/schema/UnsetTemplateProcedure.java | 10 +-
.../impl/sync/AuthOperationProcedure.java | 4 +-
.../client/request/AsyncRequestManager.java | 6 +
.../commons/exception/IoTDBRuntimeException.java | 66 ++++++++++
.../exception/UncheckedStartupException.java | 41 ++++++
44 files changed, 517 insertions(+), 285 deletions(-)
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/CnToDnRequestType.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java
similarity index 87%
rename from
iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/CnToDnRequestType.java
rename to
iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java
index 4b7ebd57dbf..19d55c170ed 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/CnToDnRequestType.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java
@@ -17,14 +17,10 @@
* under the License.
*/
-package org.apache.iotdb.confignode.client;
-
-public enum CnToDnRequestType {
+package org.apache.iotdb.confignode.client.async;
+public enum CnToDnAsyncRequestType {
// Node Maintenance
- DISABLE_DATA_NODE,
- STOP_DATA_NODE,
-
FLUSH,
MERGE,
FULL_MERGE,
@@ -33,8 +29,6 @@ public enum CnToDnRequestType {
LOAD_CONFIGURATION,
SET_SYSTEM_STATUS,
SET_CONFIGURATION,
- SHOW_CONFIGURATION,
-
SUBMIT_TEST_CONNECTION_TASK,
TEST_CONNECTION,
@@ -42,19 +36,11 @@ public enum CnToDnRequestType {
CREATE_DATA_REGION,
CREATE_SCHEMA_REGION,
DELETE_REGION,
-
- CREATE_NEW_REGION_PEER,
- ADD_REGION_PEER,
- REMOVE_REGION_PEER,
- DELETE_OLD_REGION_PEER,
RESET_PEER_LIST,
-
UPDATE_REGION_ROUTE_MAP,
CHANGE_REGION_LEADER,
// PartitionCache
- INVALIDATE_PARTITION_CACHE,
- INVALIDATE_PERMISSION_CACHE,
INVALIDATE_SCHEMA_CACHE,
INVALIDATE_LAST_CACHE,
CLEAR_CACHE,
@@ -87,15 +73,11 @@ public enum CnToDnRequestType {
CONSUMER_GROUP_PUSH_ALL_META,
CONSUMER_GROUP_PUSH_SINGLE_META,
- // CQ
- EXECUTE_CQ,
-
// TEMPLATE
UPDATE_TEMPLATE,
// Schema
SET_TTL,
- UPDATE_TTL_CACHE,
CONSTRUCT_SCHEMA_BLACK_LIST,
ROLLBACK_SCHEMA_BLACK_LIST,
@@ -113,8 +95,8 @@ public enum CnToDnRequestType {
CONSTRUCT_VIEW_SCHEMA_BLACK_LIST,
ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
- DELETE_VIEW,
+ DELETE_VIEW,
ALTER_VIEW,
// TODO Need to migrate to Node Maintenance
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java
index 6a47a7ca433..ed8e7a08aac 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java
@@ -32,7 +32,7 @@ import
org.apache.iotdb.commons.client.request.AsyncRequestContext;
import org.apache.iotdb.commons.client.request.AsyncRequestRPCHandler;
import
org.apache.iotdb.commons.client.request.DataNodeInternalServiceRequestManager;
import org.apache.iotdb.commons.client.request.TestConnectionUtils;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.commons.exception.UncheckedStartupException;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.CheckTimeSeriesExistenceRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.CountPathsUsingTemplateRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.DataNodeAsyncRequestRPCHandler;
@@ -91,9 +91,13 @@ import
org.apache.iotdb.mpp.rpc.thrift.TUpdateTriggerLocationReq;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.Arrays;
+import java.util.List;
+import java.util.stream.Collectors;
+
/** Asynchronously send RPC requests to DataNodes. See queryengine.thrift for
more details. */
public class CnToDnInternalServiceAsyncRequestManager
- extends DataNodeInternalServiceRequestManager<CnToDnRequestType> {
+ extends DataNodeInternalServiceRequestManager<CnToDnAsyncRequestType> {
private static final Logger LOGGER =
LoggerFactory.getLogger(CnToDnInternalServiceAsyncRequestManager.class);
@@ -101,272 +105,284 @@ public class CnToDnInternalServiceAsyncRequestManager
@Override
protected void initActionMapBuilder() {
actionMapBuilder.put(
- CnToDnRequestType.SET_TTL,
+ CnToDnAsyncRequestType.SET_TTL,
(req, client, handler) ->
client.setTTL((TSetTTLReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CREATE_DATA_REGION,
+ CnToDnAsyncRequestType.CREATE_DATA_REGION,
(req, client, handler) ->
client.createDataRegion(
(TCreateDataRegionReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.DELETE_REGION,
+ CnToDnAsyncRequestType.DELETE_REGION,
(req, client, handler) ->
client.deleteRegion((TConsensusGroupId) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CREATE_SCHEMA_REGION,
+ CnToDnAsyncRequestType.CREATE_SCHEMA_REGION,
(req, client, handler) ->
client.createSchemaRegion(
(TCreateSchemaRegionReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CREATE_FUNCTION,
+ CnToDnAsyncRequestType.CREATE_FUNCTION,
(req, client, handler) ->
client.createFunction(
(TCreateFunctionInstanceReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.DROP_FUNCTION,
+ CnToDnAsyncRequestType.DROP_FUNCTION,
(req, client, handler) ->
client.dropFunction(
(TDropFunctionInstanceReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CREATE_TRIGGER_INSTANCE,
+ CnToDnAsyncRequestType.CREATE_TRIGGER_INSTANCE,
(req, client, handler) ->
client.createTriggerInstance(
(TCreateTriggerInstanceReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.DROP_TRIGGER_INSTANCE,
+ CnToDnAsyncRequestType.DROP_TRIGGER_INSTANCE,
(req, client, handler) ->
client.dropTriggerInstance(
(TDropTriggerInstanceReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.ACTIVE_TRIGGER_INSTANCE,
+ CnToDnAsyncRequestType.ACTIVE_TRIGGER_INSTANCE,
(req, client, handler) ->
client.activeTriggerInstance(
(TActiveTriggerInstanceReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.INACTIVE_TRIGGER_INSTANCE,
+ CnToDnAsyncRequestType.INACTIVE_TRIGGER_INSTANCE,
(req, client, handler) ->
client.inactiveTriggerInstance(
(TInactiveTriggerInstanceReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.UPDATE_TRIGGER_LOCATION,
+ CnToDnAsyncRequestType.UPDATE_TRIGGER_LOCATION,
(req, client, handler) ->
client.updateTriggerLocation(
(TUpdateTriggerLocationReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CREATE_PIPE_PLUGIN,
+ CnToDnAsyncRequestType.CREATE_PIPE_PLUGIN,
(req, client, handler) ->
client.createPipePlugin(
(TCreatePipePluginInstanceReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.DROP_PIPE_PLUGIN,
+ CnToDnAsyncRequestType.DROP_PIPE_PLUGIN,
(req, client, handler) ->
client.dropPipePlugin(
(TDropPipePluginInstanceReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.PIPE_PUSH_ALL_META,
+ CnToDnAsyncRequestType.PIPE_PUSH_ALL_META,
(req, client, handler) ->
client.pushPipeMeta((TPushPipeMetaReq) req,
(PipePushMetaRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.PIPE_PUSH_SINGLE_META,
+ CnToDnAsyncRequestType.PIPE_PUSH_SINGLE_META,
(req, client, handler) ->
client.pushSinglePipeMeta(
(TPushSinglePipeMetaReq) req, (PipePushMetaRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.PIPE_PUSH_MULTI_META,
+ CnToDnAsyncRequestType.PIPE_PUSH_MULTI_META,
(req, client, handler) ->
client.pushMultiPipeMeta(
(TPushMultiPipeMetaReq) req, (PipePushMetaRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.TOPIC_PUSH_ALL_META,
+ CnToDnAsyncRequestType.TOPIC_PUSH_ALL_META,
(req, client, handler) ->
client.pushTopicMeta((TPushTopicMetaReq) req,
(TopicPushMetaRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.TOPIC_PUSH_SINGLE_META,
+ CnToDnAsyncRequestType.TOPIC_PUSH_SINGLE_META,
(req, client, handler) ->
client.pushSingleTopicMeta(
(TPushSingleTopicMetaReq) req, (TopicPushMetaRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.TOPIC_PUSH_MULTI_META,
+ CnToDnAsyncRequestType.TOPIC_PUSH_MULTI_META,
(req, client, handler) ->
client.pushMultiTopicMeta(
(TPushMultiTopicMetaReq) req, (TopicPushMetaRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CONSUMER_GROUP_PUSH_ALL_META,
+ CnToDnAsyncRequestType.CONSUMER_GROUP_PUSH_ALL_META,
(req, client, handler) ->
client.pushConsumerGroupMeta(
(TPushConsumerGroupMetaReq) req,
(ConsumerGroupPushMetaRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CONSUMER_GROUP_PUSH_SINGLE_META,
+ CnToDnAsyncRequestType.CONSUMER_GROUP_PUSH_SINGLE_META,
(req, client, handler) ->
client.pushSingleConsumerGroupMeta(
(TPushSingleConsumerGroupMetaReq) req,
(ConsumerGroupPushMetaRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.PIPE_HEARTBEAT,
+ CnToDnAsyncRequestType.PIPE_HEARTBEAT,
(req, client, handler) ->
client.pipeHeartbeat((TPipeHeartbeatReq) req,
(PipeHeartbeatRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.MERGE,
+ CnToDnAsyncRequestType.MERGE,
(req, client, handler) -> client.merge((DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.FULL_MERGE,
+ CnToDnAsyncRequestType.FULL_MERGE,
(req, client, handler) -> client.merge((DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.FLUSH,
+ CnToDnAsyncRequestType.FLUSH,
(req, client, handler) ->
client.flush((TFlushReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CLEAR_CACHE,
+ CnToDnAsyncRequestType.CLEAR_CACHE,
(req, client, handler) ->
client.clearCache((DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.START_REPAIR_DATA,
+ CnToDnAsyncRequestType.START_REPAIR_DATA,
(req, client, handler) ->
client.startRepairData((DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.STOP_REPAIR_DATA,
+ CnToDnAsyncRequestType.STOP_REPAIR_DATA,
(req, client, handler) ->
client.stopRepairData((DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.LOAD_CONFIGURATION,
+ CnToDnAsyncRequestType.LOAD_CONFIGURATION,
(req, client, handler) ->
client.loadConfiguration((DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.SET_SYSTEM_STATUS,
+ CnToDnAsyncRequestType.SET_SYSTEM_STATUS,
(req, client, handler) ->
client.setSystemStatus((String) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.SET_CONFIGURATION,
+ CnToDnAsyncRequestType.SET_CONFIGURATION,
(req, client, handler) ->
client.setConfiguration(
(TSetConfigurationReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.UPDATE_REGION_ROUTE_MAP,
+ CnToDnAsyncRequestType.UPDATE_REGION_ROUTE_MAP,
(req, client, handler) ->
client.updateRegionCache((TRegionRouteReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CHANGE_REGION_LEADER,
+ CnToDnAsyncRequestType.CHANGE_REGION_LEADER,
(req, client, handler) ->
client.changeRegionLeader(
(TRegionLeaderChangeReq) req, (TransferLeaderRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.CONSTRUCT_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.CONSTRUCT_SCHEMA_BLACK_LIST,
(req, client, handler) ->
client.constructSchemaBlackList(
(TConstructSchemaBlackListReq) req, (SchemaUpdateRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.ROLLBACK_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.ROLLBACK_SCHEMA_BLACK_LIST,
(req, client, handler) ->
client.rollbackSchemaBlackList(
(TRollbackSchemaBlackListReq) req, (SchemaUpdateRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.FETCH_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.FETCH_SCHEMA_BLACK_LIST,
(req, client, handler) ->
client.fetchSchemaBlackList(
(TFetchSchemaBlackListReq) req,
(FetchSchemaBlackListRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.INVALIDATE_SCHEMA_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_SCHEMA_CACHE,
(req, client, handler) ->
client.invalidateSchemaCache(
(TInvalidateCacheReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
(req, client, handler) ->
client.invalidateMatchedSchemaCache(
(TInvalidateMatchedSchemaCacheReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.INVALIDATE_LAST_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_LAST_CACHE,
(req, client, handler) ->
client.invalidateLastCache((String) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.DELETE_DATA_FOR_DELETE_SCHEMA,
+ CnToDnAsyncRequestType.DELETE_DATA_FOR_DELETE_SCHEMA,
(req, client, handler) ->
client.deleteDataForDeleteSchema(
(TDeleteDataForDeleteSchemaReq) req, (SchemaUpdateRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.DELETE_TIMESERIES,
+ CnToDnAsyncRequestType.DELETE_TIMESERIES,
(req, client, handler) ->
client.deleteTimeSeries((TDeleteTimeSeriesReq) req,
(SchemaUpdateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CONSTRUCT_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
+ CnToDnAsyncRequestType.CONSTRUCT_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
(req, client, handler) ->
client.constructSchemaBlackListWithTemplate(
(TConstructSchemaBlackListWithTemplateReq) req,
(SchemaUpdateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.ROLLBACK_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
+ CnToDnAsyncRequestType.ROLLBACK_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
(req, client, handler) ->
client.rollbackSchemaBlackListWithTemplate(
(TRollbackSchemaBlackListWithTemplateReq) req,
(SchemaUpdateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.DEACTIVATE_TEMPLATE,
+ CnToDnAsyncRequestType.DEACTIVATE_TEMPLATE,
(req, client, handler) ->
client.deactivateTemplate(
(TDeactivateTemplateReq) req, (SchemaUpdateRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.UPDATE_TEMPLATE,
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE,
(req, client, handler) ->
client.updateTemplate((TUpdateTemplateReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.COUNT_PATHS_USING_TEMPLATE,
+ CnToDnAsyncRequestType.COUNT_PATHS_USING_TEMPLATE,
(req, client, handler) ->
client.countPathsUsingTemplate(
(TCountPathsUsingTemplateReq) req,
(CountPathsUsingTemplateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CHECK_SCHEMA_REGION_USING_TEMPLATE,
+ CnToDnAsyncRequestType.CHECK_SCHEMA_REGION_USING_TEMPLATE,
(req, client, handler) ->
client.checkSchemaRegionUsingTemplate(
(TCheckSchemaRegionUsingTemplateReq) req,
(CheckSchemaRegionUsingTemplateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CHECK_TIMESERIES_EXISTENCE,
+ CnToDnAsyncRequestType.CHECK_TIMESERIES_EXISTENCE,
(req, client, handler) ->
client.checkTimeSeriesExistence(
(TCheckTimeSeriesExistenceReq) req,
(CheckTimeSeriesExistenceRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.CONSTRUCT_VIEW_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.CONSTRUCT_VIEW_SCHEMA_BLACK_LIST,
(req, client, handler) ->
client.constructViewSchemaBlackList(
(TConstructViewSchemaBlackListReq) req,
(SchemaUpdateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
(req, client, handler) ->
client.rollbackViewSchemaBlackList(
(TRollbackViewSchemaBlackListReq) req,
(SchemaUpdateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.DELETE_VIEW,
+ CnToDnAsyncRequestType.DELETE_VIEW,
(req, client, handler) ->
client.deleteViewSchema((TDeleteViewSchemaReq) req,
(SchemaUpdateRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.ALTER_VIEW,
+ CnToDnAsyncRequestType.ALTER_VIEW,
(req, client, handler) ->
client.alterView((TAlterViewReq) req, (SchemaUpdateRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.KILL_QUERY_INSTANCE,
+ CnToDnAsyncRequestType.KILL_QUERY_INSTANCE,
(req, client, handler) ->
client.killQueryInstance((String) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.SET_SPACE_QUOTA,
+ CnToDnAsyncRequestType.SET_SPACE_QUOTA,
(req, client, handler) ->
client.setSpaceQuota((TSetSpaceQuotaReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.SET_THROTTLE_QUOTA,
+ CnToDnAsyncRequestType.SET_THROTTLE_QUOTA,
(req, client, handler) ->
client.setThrottleQuota(
(TSetThrottleQuotaReq) req, (DataNodeTSStatusRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.RESET_PEER_LIST,
+ CnToDnAsyncRequestType.RESET_PEER_LIST,
(req, client, handler) ->
client.resetPeerList((TResetPeerListReq) req,
(DataNodeTSStatusRPCHandler) handler));
actionMapBuilder.put(
- CnToDnRequestType.SUBMIT_TEST_CONNECTION_TASK,
+ CnToDnAsyncRequestType.SUBMIT_TEST_CONNECTION_TASK,
(req, client, handler) ->
client.submitTestConnectionTask(
(TNodeLocations) req, (SubmitTestConnectionTaskRPCHandler)
handler));
actionMapBuilder.put(
- CnToDnRequestType.TEST_CONNECTION,
+ CnToDnAsyncRequestType.TEST_CONNECTION,
(req, client, handler) ->
client.testConnectionEmptyRPC((DataNodeTSStatusRPCHandler)
handler));
}
@Override
- protected AsyncRequestRPCHandler<?, CnToDnRequestType, TDataNodeLocation>
buildHandler(
- AsyncRequestContext<?, ?, CnToDnRequestType, TDataNodeLocation>
requestContext,
+ protected void checkActionMapCompleteness() {
+ List<CnToDnAsyncRequestType> lackList =
+ Arrays.stream(CnToDnAsyncRequestType.values())
+ .filter(type -> !actionMap.containsKey(type))
+ .collect(Collectors.toList());
+ if (!lackList.isEmpty()) {
+ throw new UncheckedStartupException(
+ String.format("These request types should be added to actionMap:
%s", lackList));
+ }
+ }
+
+ @Override
+ protected AsyncRequestRPCHandler<?, CnToDnAsyncRequestType,
TDataNodeLocation> buildHandler(
+ AsyncRequestContext<?, ?, CnToDnAsyncRequestType, TDataNodeLocation>
requestContext,
int requestId,
TDataNodeLocation targetNode) {
return DataNodeAsyncRequestRPCHandler.buildHandler(requestContext,
requestId, targetNode);
@@ -374,8 +390,8 @@ public class CnToDnInternalServiceAsyncRequestManager
@Override
protected void adjustClientTimeoutIfNecessary(
- CnToDnRequestType cnToDnRequestType, AsyncDataNodeInternalServiceClient
client) {
- if
(CnToDnRequestType.SUBMIT_TEST_CONNECTION_TASK.equals(cnToDnRequestType)) {
+ CnToDnAsyncRequestType CnToDnAsyncRequestType,
AsyncDataNodeInternalServiceClient client) {
+ if
(CnToDnAsyncRequestType.SUBMIT_TEST_CONNECTION_TASK.equals(CnToDnAsyncRequestType))
{
client.setTimeoutTemporarily(TestConnectionUtils.calculateCnLeaderToAllDnMaxTime());
}
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/DataNodeAsyncRequestContext.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/DataNodeAsyncRequestContext.java
index 2b813c081d7..50797d603a3 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/DataNodeAsyncRequestContext.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/DataNodeAsyncRequestContext.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.client.async.handlers;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.commons.client.request.AsyncRequestContext;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import java.util.Map;
@@ -32,19 +32,21 @@ import java.util.Map;
* @param <R> ClassName of RPC response
*/
public class DataNodeAsyncRequestContext<Q, R>
- extends AsyncRequestContext<Q, R, CnToDnRequestType, TDataNodeLocation> {
+ extends AsyncRequestContext<Q, R, CnToDnAsyncRequestType,
TDataNodeLocation> {
- public DataNodeAsyncRequestContext(CnToDnRequestType requestType) {
+ public DataNodeAsyncRequestContext(CnToDnAsyncRequestType requestType) {
super(requestType);
}
public DataNodeAsyncRequestContext(
- CnToDnRequestType requestType, Map<Integer, TDataNodeLocation>
dataNodeLocationMap) {
+ CnToDnAsyncRequestType requestType, Map<Integer, TDataNodeLocation>
dataNodeLocationMap) {
super(requestType, dataNodeLocationMap);
}
public DataNodeAsyncRequestContext(
- CnToDnRequestType requestType, Q q, Map<Integer, TDataNodeLocation>
dataNodeLocationMap) {
+ CnToDnAsyncRequestType requestType,
+ Q q,
+ Map<Integer, TDataNodeLocation> dataNodeLocationMap) {
super(requestType, q, dataNodeLocationMap);
}
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CheckTimeSeriesExistenceRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CheckTimeSeriesExistenceRPCHandler.java
index b12ba289f1c..3a735691a0e 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CheckTimeSeriesExistenceRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CheckTimeSeriesExistenceRPCHandler.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.mpp.rpc.thrift.TCheckTimeSeriesExistenceResp;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -39,7 +39,7 @@ public class CheckTimeSeriesExistenceRPCHandler
LoggerFactory.getLogger(CheckTimeSeriesExistenceRPCHandler.class);
public CheckTimeSeriesExistenceRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CountPathsUsingTemplateRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CountPathsUsingTemplateRPCHandler.java
index f2ed0cc425c..b27c74bb41d 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CountPathsUsingTemplateRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/CountPathsUsingTemplateRPCHandler.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.mpp.rpc.thrift.TCountPathsUsingTemplateResp;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -39,7 +39,7 @@ public class CountPathsUsingTemplateRPCHandler
LoggerFactory.getLogger(CountPathsUsingTemplateRPCHandler.class);
public CountPathsUsingTemplateRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeAsyncRequestRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeAsyncRequestRPCHandler.java
index 19be87ef068..fe09c02e92b 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeAsyncRequestRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeAsyncRequestRPCHandler.java
@@ -24,7 +24,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.common.rpc.thrift.TTestConnectionResp;
import org.apache.iotdb.commons.client.request.AsyncRequestContext;
import org.apache.iotdb.commons.client.request.AsyncRequestRPCHandler;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.subscription.CheckSchemaRegionUsingTemplateRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.subscription.ConsumerGroupPushMetaRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.subscription.TopicPushMetaRPCHandler;
@@ -42,10 +42,10 @@ import java.util.Map;
import java.util.concurrent.CountDownLatch;
public abstract class DataNodeAsyncRequestRPCHandler<Response>
- extends AsyncRequestRPCHandler<Response, CnToDnRequestType,
TDataNodeLocation> {
+ extends AsyncRequestRPCHandler<Response, CnToDnAsyncRequestType,
TDataNodeLocation> {
protected DataNodeAsyncRequestRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
@@ -70,10 +70,10 @@ public abstract class
DataNodeAsyncRequestRPCHandler<Response>
}
public static DataNodeAsyncRequestRPCHandler<?> buildHandler(
- AsyncRequestContext<?, ?, CnToDnRequestType, TDataNodeLocation> context,
+ AsyncRequestContext<?, ?, CnToDnAsyncRequestType, TDataNodeLocation>
context,
int requestId,
TDataNodeLocation targetDataNode) {
- CnToDnRequestType requestType = context.getRequestType();
+ CnToDnAsyncRequestType requestType = context.getRequestType();
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
context.getNodeLocationMap();
Map<Integer, ?> responseMap = context.getResponseMap();
CountDownLatch countDownLatch = context.getCountDownLatch();
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
index 19d451eb671..7c93f363dd4 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -37,7 +37,7 @@ public class DataNodeTSStatusRPCHandler extends
DataNodeAsyncRequestRPCHandler<T
private static final Logger LOGGER =
LoggerFactory.getLogger(DataNodeTSStatusRPCHandler.class);
public DataNodeTSStatusRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/FetchSchemaBlackListRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/FetchSchemaBlackListRPCHandler.java
index 45c659298ca..693017ec02d 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/FetchSchemaBlackListRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/FetchSchemaBlackListRPCHandler.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.mpp.rpc.thrift.TFetchSchemaBlackListResp;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -39,7 +39,7 @@ public class FetchSchemaBlackListRPCHandler
LoggerFactory.getLogger(FetchSchemaBlackListRPCHandler.class);
public FetchSchemaBlackListRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipeHeartbeatRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipeHeartbeatRPCHandler.java
index ec4968cfa06..569424afbff 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipeHeartbeatRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipeHeartbeatRPCHandler.java
@@ -20,7 +20,7 @@
package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.mpp.rpc.thrift.TPipeHeartbeatResp;
import org.slf4j.Logger;
@@ -34,7 +34,7 @@ public class PipeHeartbeatRPCHandler extends
DataNodeAsyncRequestRPCHandler<TPip
private static final Logger LOGGER =
LoggerFactory.getLogger(PipeHeartbeatRPCHandler.class);
public PipeHeartbeatRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipePushMetaRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipePushMetaRPCHandler.java
index 3517bbeb93e..9ef80e9f847 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipePushMetaRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/PipePushMetaRPCHandler.java
@@ -20,7 +20,7 @@
package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaResp;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -35,7 +35,7 @@ public class PipePushMetaRPCHandler extends
DataNodeAsyncRequestRPCHandler<TPush
private static final Logger LOGGER =
LoggerFactory.getLogger(PipePushMetaRPCHandler.class);
public PipePushMetaRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
index db8458a948c..dc2796a232e 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -36,7 +36,7 @@ public class SchemaUpdateRPCHandler extends
DataNodeTSStatusRPCHandler {
private static final Logger LOGGER =
LoggerFactory.getLogger(SchemaUpdateRPCHandler.class);
public SchemaUpdateRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SubmitTestConnectionTaskRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SubmitTestConnectionTaskRPCHandler.java
index 4abf0a0eca1..f3c58892cbe 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SubmitTestConnectionTaskRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SubmitTestConnectionTaskRPCHandler.java
@@ -22,7 +22,7 @@ package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.common.rpc.thrift.TTestConnectionResp;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.rpc.TSStatusCode;
import org.slf4j.Logger;
@@ -38,7 +38,7 @@ public class SubmitTestConnectionTaskRPCHandler
LoggerFactory.getLogger(SubmitTestConnectionTaskRPCHandler.class);
public SubmitTestConnectionTaskRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/TransferLeaderRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/TransferLeaderRPCHandler.java
index 352cc0694e2..8bfe0eb4755 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/TransferLeaderRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/TransferLeaderRPCHandler.java
@@ -20,7 +20,7 @@
package org.apache.iotdb.confignode.client.async.handlers.rpc;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.mpp.rpc.thrift.TRegionLeaderChangeResp;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -37,7 +37,7 @@ public class TransferLeaderRPCHandler
private static final Logger LOGGER =
LoggerFactory.getLogger(TransferLeaderRPCHandler.class);
public TransferLeaderRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/CheckSchemaRegionUsingTemplateRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/CheckSchemaRegionUsingTemplateRPCHandler.java
index 14898dcdc6c..249e8b51767 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/CheckSchemaRegionUsingTemplateRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/CheckSchemaRegionUsingTemplateRPCHandler.java
@@ -21,7 +21,7 @@ package
org.apache.iotdb.confignode.client.async.handlers.rpc.subscription;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.DataNodeAsyncRequestRPCHandler;
import org.apache.iotdb.mpp.rpc.thrift.TCheckSchemaRegionUsingTemplateResp;
import org.apache.iotdb.rpc.RpcUtils;
@@ -40,7 +40,7 @@ public class CheckSchemaRegionUsingTemplateRPCHandler
LoggerFactory.getLogger(CheckSchemaRegionUsingTemplateRPCHandler.class);
public CheckSchemaRegionUsingTemplateRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/ConsumerGroupPushMetaRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/ConsumerGroupPushMetaRPCHandler.java
index ee3c11eeb42..2938d4f85b7 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/ConsumerGroupPushMetaRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/ConsumerGroupPushMetaRPCHandler.java
@@ -20,7 +20,7 @@
package org.apache.iotdb.confignode.client.async.handlers.rpc.subscription;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.DataNodeAsyncRequestRPCHandler;
import org.apache.iotdb.mpp.rpc.thrift.TPushConsumerGroupMetaResp;
import org.apache.iotdb.rpc.RpcUtils;
@@ -38,7 +38,7 @@ public class ConsumerGroupPushMetaRPCHandler
LoggerFactory.getLogger(ConsumerGroupPushMetaRPCHandler.class);
public ConsumerGroupPushMetaRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/TopicPushMetaRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/TopicPushMetaRPCHandler.java
index cf8451feaac..91ffdd7232b 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/TopicPushMetaRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/subscription/TopicPushMetaRPCHandler.java
@@ -20,7 +20,7 @@
package org.apache.iotdb.confignode.client.async.handlers.rpc.subscription;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.DataNodeAsyncRequestRPCHandler;
import org.apache.iotdb.mpp.rpc.thrift.TPushTopicMetaResp;
import org.apache.iotdb.rpc.RpcUtils;
@@ -37,7 +37,7 @@ public class TopicPushMetaRPCHandler extends
DataNodeAsyncRequestRPCHandler<TPus
private static final Logger LOGGER =
LoggerFactory.getLogger(TopicPushMetaRPCHandler.class);
public TopicPushMetaRPCHandler(
- CnToDnRequestType requestType,
+ CnToDnAsyncRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
Map<Integer, TDataNodeLocation> dataNodeLocationMap,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/CnToDnSyncRequestType.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/CnToDnSyncRequestType.java
new file mode 100644
index 00000000000..32d1d854c6c
--- /dev/null
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/CnToDnSyncRequestType.java
@@ -0,0 +1,49 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.confignode.client.sync;
+
+public enum CnToDnSyncRequestType {
+ // Node Maintenance
+ DISABLE_DATANODE,
+ STOP_DATA_NODE,
+ SET_SYSTEM_STATUS,
+ SHOW_CONFIGURATION,
+
+ // Region Maintenance
+ CREATE_DATA_REGION,
+ CREATE_SCHEMA_REGION,
+ DELETE_REGION,
+ CREATE_NEW_REGION_PEER,
+ ADD_REGION_PEER,
+ REMOVE_REGION_PEER,
+ DELETE_OLD_REGION_PEER,
+ RESET_PEER_LIST,
+
+ // PartitionCache
+ INVALIDATE_PARTITION_CACHE,
+ INVALIDATE_PERMISSION_CACHE,
+ INVALIDATE_SCHEMA_CACHE,
+
+ // Template
+ UPDATE_TEMPLATE,
+
+ // Schema
+ KILL_QUERY_INSTANCE,
+}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/SyncDataNodeClientPool.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/SyncDataNodeClientPool.java
index 84e609ffb38..d618e843250 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/SyncDataNodeClientPool.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/sync/SyncDataNodeClientPool.java
@@ -27,7 +27,7 @@ import org.apache.iotdb.commons.client.ClientPoolFactory;
import org.apache.iotdb.commons.client.IClientManager;
import org.apache.iotdb.commons.client.exception.ClientManagerException;
import org.apache.iotdb.commons.client.sync.SyncDataNodeInternalServiceClient;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.commons.exception.UncheckedStartupException;
import org.apache.iotdb.mpp.rpc.thrift.TCreateDataRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreatePeerReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateSchemaRegionReq;
@@ -39,14 +39,19 @@ import
org.apache.iotdb.mpp.rpc.thrift.TRegionLeaderChangeReq;
import org.apache.iotdb.mpp.rpc.thrift.TRegionLeaderChangeResp;
import org.apache.iotdb.mpp.rpc.thrift.TResetPeerListReq;
import org.apache.iotdb.mpp.rpc.thrift.TUpdateTemplateReq;
-import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
+import com.google.common.collect.ImmutableMap;
+import org.apache.ratis.util.function.CheckedBiFunction;
import org.apache.thrift.TException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Objects;
import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
/** Synchronously send RPC requests to DataNodes. See queryengine.thrift for
more details. */
public class SyncDataNodeClientPool {
@@ -57,20 +62,95 @@ public class SyncDataNodeClientPool {
private final IClientManager<TEndPoint, SyncDataNodeInternalServiceClient>
clientManager;
+ protected ImmutableMap<
+ CnToDnSyncRequestType,
+ CheckedBiFunction<Object, SyncDataNodeInternalServiceClient, Object,
Exception>>
+ actionMap;
+
private SyncDataNodeClientPool() {
clientManager =
new IClientManager.Factory<TEndPoint,
SyncDataNodeInternalServiceClient>()
.createClientManager(
new
ClientPoolFactory.SyncDataNodeInternalServiceClientPoolFactory());
+ buildActionMap();
+ checkActionMapCompleteness();
+ }
+
+ private void buildActionMap() {
+ ImmutableMap.Builder<
+ CnToDnSyncRequestType,
+ CheckedBiFunction<Object, SyncDataNodeInternalServiceClient,
Object, Exception>>
+ actionMapBuilder = ImmutableMap.builder();
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.INVALIDATE_PARTITION_CACHE,
+ (req, client) -> client.invalidatePartitionCache((TInvalidateCacheReq)
req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.INVALIDATE_SCHEMA_CACHE,
+ (req, client) -> client.invalidateSchemaCache((TInvalidateCacheReq)
req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.CREATE_SCHEMA_REGION,
+ (req, client) -> client.createSchemaRegion((TCreateSchemaRegionReq)
req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.CREATE_DATA_REGION,
+ (req, client) -> client.createDataRegion((TCreateDataRegionReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.DELETE_REGION,
+ (req, client) -> client.deleteRegion((TConsensusGroupId) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.INVALIDATE_PERMISSION_CACHE,
+ (req, client) ->
client.invalidatePermissionCache((TInvalidatePermissionCacheReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.DISABLE_DATANODE,
+ (req, client) -> client.disableDataNode((TDisableDataNodeReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.STOP_DATA_NODE, (req, client) ->
client.stopDataNode());
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.SET_SYSTEM_STATUS,
+ (req, client) -> client.setSystemStatus((String) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.KILL_QUERY_INSTANCE,
+ (req, client) -> client.killQueryInstance((String) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.UPDATE_TEMPLATE,
+ (req, client) -> client.updateTemplate((TUpdateTemplateReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.CREATE_NEW_REGION_PEER,
+ (req, client) -> client.createNewRegionPeer((TCreatePeerReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.ADD_REGION_PEER,
+ (req, client) -> client.addRegionPeer((TMaintainPeerReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.REMOVE_REGION_PEER,
+ (req, client) -> client.removeRegionPeer((TMaintainPeerReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.DELETE_OLD_REGION_PEER,
+ (req, client) -> client.deleteOldRegionPeer((TMaintainPeerReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.RESET_PEER_LIST,
+ (req, client) -> client.resetPeerList((TResetPeerListReq) req));
+ actionMapBuilder.put(
+ CnToDnSyncRequestType.SHOW_CONFIGURATION, (req, client) ->
client.showConfiguration());
+ actionMap = actionMapBuilder.build();
+ }
+
+ private void checkActionMapCompleteness() {
+ List<CnToDnSyncRequestType> lackList =
+ Arrays.stream(CnToDnSyncRequestType.values())
+ .filter(type -> !actionMap.containsKey(type))
+ .collect(Collectors.toList());
+ if (!lackList.isEmpty()) {
+ throw new UncheckedStartupException(
+ String.format("These request types should be added to actionMap:
%s", lackList));
+ }
}
public Object sendSyncRequestToDataNodeWithRetry(
- TEndPoint endPoint, Object req, CnToDnRequestType requestType) {
+ TEndPoint endPoint, Object req, CnToDnSyncRequestType requestType) {
Throwable lastException = new TException();
for (int retry = 0; retry < DEFAULT_RETRY_NUM; retry++) {
try (SyncDataNodeInternalServiceClient client =
clientManager.borrowClient(endPoint)) {
return executeSyncRequest(requestType, client, req);
- } catch (ClientManagerException | TException e) {
+ } catch (Exception e) {
lastException = e;
if (retry != DEFAULT_RETRY_NUM - 1) {
LOGGER.warn("{} failed on DataNode {}, retrying {}...", requestType,
endPoint, retry + 1);
@@ -84,12 +164,12 @@ public class SyncDataNodeClientPool {
}
public Object sendSyncRequestToDataNodeWithGivenRetry(
- TEndPoint endPoint, Object req, CnToDnRequestType requestType, int
retryNum) {
+ TEndPoint endPoint, Object req, CnToDnSyncRequestType requestType, int
retryNum) {
Throwable lastException = new TException();
for (int retry = 0; retry < retryNum; retry++) {
try (SyncDataNodeInternalServiceClient client =
clientManager.borrowClient(endPoint)) {
return executeSyncRequest(requestType, client, req);
- } catch (ClientManagerException | TException e) {
+ } catch (Exception e) {
lastException = e;
if (retry != retryNum - 1) {
LOGGER.warn("{} failed on DataNode {}, retrying {}...", requestType,
endPoint, retry + 1);
@@ -103,47 +183,9 @@ public class SyncDataNodeClientPool {
}
private Object executeSyncRequest(
- CnToDnRequestType requestType, SyncDataNodeInternalServiceClient client,
Object req)
- throws TException {
- switch (requestType) {
- case INVALIDATE_PARTITION_CACHE:
- return client.invalidatePartitionCache((TInvalidateCacheReq) req);
- case INVALIDATE_SCHEMA_CACHE:
- return client.invalidateSchemaCache((TInvalidateCacheReq) req);
- case CREATE_SCHEMA_REGION:
- return client.createSchemaRegion((TCreateSchemaRegionReq) req);
- case CREATE_DATA_REGION:
- return client.createDataRegion((TCreateDataRegionReq) req);
- case DELETE_REGION:
- return client.deleteRegion((TConsensusGroupId) req);
- case INVALIDATE_PERMISSION_CACHE:
- return
client.invalidatePermissionCache((TInvalidatePermissionCacheReq) req);
- case DISABLE_DATA_NODE:
- return client.disableDataNode((TDisableDataNodeReq) req);
- case STOP_DATA_NODE:
- return client.stopDataNode();
- case SET_SYSTEM_STATUS:
- return client.setSystemStatus((String) req);
- case KILL_QUERY_INSTANCE:
- return client.killQueryInstance((String) req);
- case UPDATE_TEMPLATE:
- return client.updateTemplate((TUpdateTemplateReq) req);
- case CREATE_NEW_REGION_PEER:
- return client.createNewRegionPeer((TCreatePeerReq) req);
- case ADD_REGION_PEER:
- return client.addRegionPeer((TMaintainPeerReq) req);
- case REMOVE_REGION_PEER:
- return client.removeRegionPeer((TMaintainPeerReq) req);
- case DELETE_OLD_REGION_PEER:
- return client.deleteOldRegionPeer((TMaintainPeerReq) req);
- case RESET_PEER_LIST:
- return client.resetPeerList((TResetPeerListReq) req);
- case SHOW_CONFIGURATION:
- return client.showConfiguration();
- default:
- return RpcUtils.getStatus(
- TSStatusCode.EXECUTE_STATEMENT_ERROR, "Unknown request type: " +
requestType);
- }
+ CnToDnSyncRequestType requestType, SyncDataNodeInternalServiceClient
client, Object req)
+ throws Exception {
+ return Objects.requireNonNull(actionMap.get(requestType)).apply(req,
client);
}
private void doRetryWait(int retryNum) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
index ced2964ec5c..3400b5316fc 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
@@ -25,6 +25,8 @@ import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.exception.ConfigurationException;
import org.apache.iotdb.commons.exception.StartupException;
import org.apache.iotdb.commons.service.StartupChecks;
+import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
+import org.apache.iotdb.confignode.client.sync.SyncDataNodeClientPool;
import
org.apache.iotdb.confignode.manager.load.balancer.router.leader.AbstractLeaderBalancer;
import
org.apache.iotdb.confignode.manager.load.balancer.router.priority.IPriorityBalancer;
import org.apache.iotdb.consensus.ConsensusFactory;
@@ -73,6 +75,7 @@ public class ConfigNodeStartupCheck extends StartupChecks {
verify();
checkGlobalConfig();
createDirsIfNecessary();
+ checkRequestManager();
if (SystemPropertiesUtils.isRestarted()) {
/* Always restore ConfigNodeId first */
CONF.setConfigNodeId(SystemPropertiesUtils.loadConfigNodeIdWhenRestarted());
@@ -223,4 +226,10 @@ public class ConfigNodeStartupCheck extends StartupChecks {
}
}
}
+
+ // The checks are in the initialization process of the RequestManager object.
+ private void checkRequestManager() {
+ SyncDataNodeClientPool.getInstance();
+ CnToDnInternalServiceAsyncRequestManager.getInstance();
+ }
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterManager.java
index 8dba6addffd..e7d0dea48f8 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterManager.java
@@ -33,8 +33,8 @@ import
org.apache.iotdb.common.rpc.thrift.TTestConnectionResult;
import org.apache.iotdb.commons.client.request.AsyncRequestContext;
import org.apache.iotdb.commons.client.request.TestConnectionUtils;
import org.apache.iotdb.confignode.client.CnToCnNodeRequestType;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
import
org.apache.iotdb.confignode.client.async.CnToCnInternalServiceAsyncRequestManager;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.ConfigNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
@@ -155,7 +155,7 @@ public class ClusterManager {
.collect(Collectors.toMap(TDataNodeLocation::getDataNodeId,
location -> location));
DataNodeAsyncRequestContext<TNodeLocations, TTestConnectionResp>
dataNodeAsyncRequestContext =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.SUBMIT_TEST_CONNECTION_TASK, nodeLocations,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.SUBMIT_TEST_CONNECTION_TASK, nodeLocations,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequest(dataNodeAsyncRequestContext);
Map<Integer, TDataNodeLocation> anotherDataNodeLocationMap =
@@ -226,8 +226,9 @@ public class ClusterManager {
TDataNodeLocation::getDataNodeId,
TDataNodeLocation::getInternalEndPoint,
TServiceType.DataNodeInternalService,
- CnToDnRequestType.TEST_CONNECTION,
- (AsyncRequestContext<Object, TSStatus, CnToDnRequestType,
TDataNodeLocation> handler) ->
+ CnToDnAsyncRequestType.TEST_CONNECTION,
+ (AsyncRequestContext<Object, TSStatus, CnToDnAsyncRequestType,
TDataNodeLocation>
+ handler) ->
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequest(handler));
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterQuotaManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterQuotaManager.java
index b0150b349d8..cbfa1bc2fb2 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterQuotaManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ClusterQuotaManager.java
@@ -26,7 +26,7 @@ import
org.apache.iotdb.common.rpc.thrift.TSetThrottleQuotaReq;
import org.apache.iotdb.common.rpc.thrift.TSpaceQuota;
import org.apache.iotdb.common.rpc.thrift.TThrottleQuota;
import org.apache.iotdb.commons.conf.IoTDBConstant;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.write.quota.SetSpaceQuotaPlan;
@@ -89,7 +89,7 @@ public class ClusterQuotaManager {
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TSetSpaceQuotaReq, TSStatus> clientHandler
=
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.SET_SPACE_QUOTA, req, dataNodeLocationMap);
+ CnToDnAsyncRequestType.SET_SPACE_QUOTA, req,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestWithRetry(clientHandler);
return
RpcUtils.squashResponseStatusList(clientHandler.getResponseList());
@@ -196,7 +196,7 @@ public class ClusterQuotaManager {
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TSetThrottleQuotaReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.SET_THROTTLE_QUOTA, req,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.SET_THROTTLE_QUOTA, req,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestWithRetry(clientHandler);
return
RpcUtils.squashResponseStatusList(clientHandler.getResponseList());
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/TriggerManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/TriggerManager.java
index b1a96d029af..5f64c496312 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/TriggerManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/TriggerManager.java
@@ -24,7 +24,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.trigger.TriggerInformation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.read.trigger.GetTransferringTriggersPlan;
@@ -250,7 +250,7 @@ public class TriggerManager {
DataNodeAsyncRequestContext<TUpdateTriggerLocationReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TRIGGER_LOCATION, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TRIGGER_LOCATION, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/UDFManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/UDFManager.java
index ad86879d828..00ed020a14e 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/UDFManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/UDFManager.java
@@ -23,7 +23,7 @@ import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.udf.UDFInformation;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
@@ -129,7 +129,7 @@ public class UDFManager {
new
TCreateFunctionInstanceReq(udfInformation.serialize()).setJarFile(jarFile);
DataNodeAsyncRequestContext<TCreateFunctionInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CREATE_FUNCTION, req, dataNodeLocationMap);
+ CnToDnAsyncRequestType.CREATE_FUNCTION, req, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -163,7 +163,7 @@ public class UDFManager {
DataNodeAsyncRequestContext<TDropFunctionInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.DROP_FUNCTION, request, dataNodeLocationMap);
+ CnToDnAsyncRequestType.DROP_FUNCTION, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RouteBalancer.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RouteBalancer.java
index ad52faa6249..350aa783360 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RouteBalancer.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RouteBalancer.java
@@ -25,7 +25,7 @@ import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.cluster.NodeStatus;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
@@ -170,7 +170,7 @@ public class RouteBalancer implements
IClusterStatusSubscriber {
long currentTime = System.nanoTime();
AtomicInteger requestId = new AtomicInteger(0);
DataNodeAsyncRequestContext<TRegionLeaderChangeReq,
TRegionLeaderChangeResp> clientHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnRequestType.CHANGE_REGION_LEADER);
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.CHANGE_REGION_LEADER);
Map<TConsensusGroupId, ConsensusGroupHeartbeatSample> successTransferMap =
new TreeMap<>();
optimalLeaderMap.forEach(
(regionGroupId, newLeaderId) -> {
@@ -255,7 +255,7 @@ public class RouteBalancer implements
IClusterStatusSubscriber {
private void invalidateSchemaCacheOfOldLeaders(
Map<TConsensusGroupId, Integer> oldLeaderMap, Set<TConsensusGroupId>
successTransferSet) {
DataNodeAsyncRequestContext<String, TSStatus>
invalidateSchemaCacheRequestHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnRequestType.INVALIDATE_LAST_CACHE);
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.INVALIDATE_LAST_CACHE);
AtomicInteger requestIndex = new AtomicInteger(0);
oldLeaderMap.entrySet().stream()
.filter(entry -> TConsensusGroupType.DataRegion ==
entry.getKey().getType())
@@ -344,7 +344,7 @@ public class RouteBalancer implements
IClusterStatusSubscriber {
Map<TConsensusGroupId, TRegionReplicaSet> tmpPriorityMap =
getRegionPriorityMap();
DataNodeAsyncRequestContext<TRegionRouteReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_REGION_ROUTE_MAP,
+ CnToDnAsyncRequestType.UPDATE_REGION_ROUTE_MAP,
new TRegionRouteReq(broadcastTime, tmpPriorityMap),
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/node/NodeManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/node/NodeManager.java
index f7cf5a78e5d..ada9262de54 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/node/NodeManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/node/NodeManager.java
@@ -37,9 +37,10 @@ import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.consensus.ConsensusGroupId;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.confignode.client.CnToCnNodeRequestType;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
+import org.apache.iotdb.confignode.client.sync.CnToDnSyncRequestType;
import org.apache.iotdb.confignode.client.sync.SyncConfigNodeClientPool;
import org.apache.iotdb.confignode.client.sync.SyncDataNodeClientPool;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
@@ -857,7 +858,7 @@ public class NodeManager {
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<Object, TSStatus> clientHandler =
- new DataNodeAsyncRequestContext<>(CnToDnRequestType.MERGE,
dataNodeLocationMap);
+ new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.MERGE,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -866,7 +867,7 @@ public class NodeManager {
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TFlushReq, TSStatus> clientHandler =
- new DataNodeAsyncRequestContext<>(CnToDnRequestType.FLUSH, req,
dataNodeLocationMap);
+ new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.FLUSH, req,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -875,7 +876,7 @@ public class NodeManager {
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<Object, TSStatus> clientHandler =
- new DataNodeAsyncRequestContext<>(CnToDnRequestType.CLEAR_CACHE,
dataNodeLocationMap);
+ new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.CLEAR_CACHE,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -896,7 +897,7 @@ public class NodeManager {
if (!targetDataNodes.isEmpty()) {
DataNodeAsyncRequestContext<Object, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.SET_CONFIGURATION, req, dataNodeLocationMap);
+ CnToDnAsyncRequestType.SET_CONFIGURATION, req,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestWithRetry(clientHandler);
responseList.addAll(clientHandler.getResponseList());
@@ -934,7 +935,8 @@ public class NodeManager {
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<Object, TSStatus> clientHandler =
- new DataNodeAsyncRequestContext<>(CnToDnRequestType.START_REPAIR_DATA,
dataNodeLocationMap);
+ new DataNodeAsyncRequestContext<>(
+ CnToDnAsyncRequestType.START_REPAIR_DATA, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -943,7 +945,8 @@ public class NodeManager {
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<Object, TSStatus> clientHandler =
- new DataNodeAsyncRequestContext<>(CnToDnRequestType.STOP_REPAIR_DATA,
dataNodeLocationMap);
+ new DataNodeAsyncRequestContext<>(
+ CnToDnAsyncRequestType.STOP_REPAIR_DATA, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -953,7 +956,7 @@ public class NodeManager {
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<Object, TSStatus> dataNodeRequestContext =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.LOAD_CONFIGURATION, dataNodeLocationMap);
+ CnToDnAsyncRequestType.LOAD_CONFIGURATION, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestWithRetry(dataNodeRequestContext);
return dataNodeRequestContext.getResponseList();
@@ -972,7 +975,7 @@ public class NodeManager {
.sendSyncRequestToDataNodeWithRetry(
dataNodeLocation.getInternalEndPoint(),
null,
- CnToDnRequestType.SHOW_CONFIGURATION);
+ CnToDnSyncRequestType.SHOW_CONFIGURATION);
}
// other config node
@@ -997,7 +1000,7 @@ public class NodeManager {
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<String, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.SET_SYSTEM_STATUS, status, dataNodeLocationMap);
+ CnToDnAsyncRequestType.SET_SYSTEM_STATUS, status,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -1008,7 +1011,7 @@ public class NodeManager {
.sendSyncRequestToDataNodeWithRetry(
setDataNodeStatusReq.getTargetDataNode().getInternalEndPoint(),
setDataNodeStatusReq.getStatus(),
- CnToDnRequestType.SET_SYSTEM_STATUS);
+ CnToDnSyncRequestType.SET_SYSTEM_STATUS);
}
/**
@@ -1031,7 +1034,7 @@ public class NodeManager {
configManager.getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<String, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.KILL_QUERY_INSTANCE, dataNodeLocationMap);
+ CnToDnAsyncRequestType.KILL_QUERY_INSTANCE, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return RpcUtils.squashResponseStatusList(clientHandler.getResponseList());
}
@@ -1047,7 +1050,7 @@ public class NodeManager {
.sendSyncRequestToDataNodeWithRetry(
dataNodeLocation.getInternalEndPoint(),
queryId,
- CnToDnRequestType.KILL_QUERY_INSTANCE);
+ CnToDnSyncRequestType.KILL_QUERY_INSTANCE);
}
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
index 642dd64e0aa..603e8829456 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
@@ -34,7 +34,7 @@ import org.apache.iotdb.commons.partition.DataPartitionTable;
import org.apache.iotdb.commons.partition.SchemaPartitionTable;
import org.apache.iotdb.commons.partition.executor.SeriesPartitionExecutor;
import org.apache.iotdb.commons.path.PartialPath;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
@@ -1260,7 +1260,7 @@ public class PartitionManager {
DataNodeAsyncRequestContext<TCreateSchemaRegionReq,
TSStatus>
createSchemaRegionHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CREATE_SCHEMA_REGION);
+
CnToDnAsyncRequestType.CREATE_SCHEMA_REGION);
for (RegionMaintainTask regionMaintainTask :
selectedRegionMaintainTask) {
RegionCreateTask schemaRegionCreateTask =
(RegionCreateTask) regionMaintainTask;
@@ -1296,7 +1296,7 @@ public class PartitionManager {
DataNodeAsyncRequestContext<TCreateDataRegionReq,
TSStatus>
createDataRegionHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CREATE_DATA_REGION);
+
CnToDnAsyncRequestType.CREATE_DATA_REGION);
for (RegionMaintainTask regionMaintainTask :
selectedRegionMaintainTask) {
RegionCreateTask dataRegionCreateTask =
(RegionCreateTask) regionMaintainTask;
@@ -1332,7 +1332,7 @@ public class PartitionManager {
case DELETE:
// delete region
DataNodeAsyncRequestContext<TConsensusGroupId, TSStatus>
deleteRegionHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnRequestType.DELETE_REGION);
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION);
Map<Integer, TConsensusGroupId> regionIdMap = new
HashMap<>();
for (RegionMaintainTask regionMaintainTask :
selectedRegionMaintainTask) {
RegionDeleteTask regionDeleteTask = (RegionDeleteTask)
regionMaintainTask;
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
index b21fbc815f9..3533b40158d 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
@@ -24,7 +24,7 @@ import
org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
import org.apache.iotdb.commons.concurrent.ThreadName;
import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
@@ -96,7 +96,7 @@ public class PipeHeartbeatScheduler {
final DataNodeAsyncRequestContext<TPipeHeartbeatReq, TPipeHeartbeatResp>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.PIPE_HEARTBEAT, request, dataNodeLocationMap);
+ CnToDnAsyncRequestType.PIPE_HEARTBEAT, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
index ebb7a264e04..25b65e606cd 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
@@ -30,7 +30,7 @@ import org.apache.iotdb.commons.schema.SchemaConstant;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.commons.utils.PathUtils;
import org.apache.iotdb.commons.utils.StatusUtils;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
@@ -1037,7 +1037,7 @@ public class ClusterSchemaManager {
DataNodeAsyncRequestContext<TUpdateTemplateReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TEMPLATE, updateTemplateReq,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE, updateTemplateReq,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
Map<Integer, TSStatus> statusMap = clientHandler.getResponseMap();
for (Map.Entry<Integer, TSStatus> entry : statusMap.entrySet()) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
index 23811a86fd0..774c01d2937 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
@@ -33,9 +33,10 @@ import
org.apache.iotdb.commons.pipe.agent.plugin.meta.PipePluginMeta;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
import org.apache.iotdb.commons.trigger.TriggerInformation;
import org.apache.iotdb.confignode.client.CnToCnNodeRequestType;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
+import org.apache.iotdb.confignode.client.sync.CnToDnSyncRequestType;
import org.apache.iotdb.confignode.client.sync.SyncConfigNodeClientPool;
import org.apache.iotdb.confignode.client.sync.SyncDataNodeClientPool;
import
org.apache.iotdb.confignode.consensus.request.write.confignode.RemoveConfigNodePlan;
@@ -186,7 +187,7 @@ public class ConfigNodeProcedureEnv {
.sendSyncRequestToDataNodeWithRetry(
dataNodeConfiguration.getLocation().getInternalEndPoint(),
invalidateCacheReq,
- CnToDnRequestType.INVALIDATE_PARTITION_CACHE);
+ CnToDnSyncRequestType.INVALIDATE_PARTITION_CACHE);
final TSStatus invalidateSchemaStatus =
(TSStatus)
@@ -194,7 +195,7 @@ public class ConfigNodeProcedureEnv {
.sendSyncRequestToDataNodeWithRetry(
dataNodeConfiguration.getLocation().getInternalEndPoint(),
invalidateCacheReq,
- CnToDnRequestType.INVALIDATE_SCHEMA_CACHE);
+ CnToDnSyncRequestType.INVALIDATE_SCHEMA_CACHE);
if (!verifySucceed(invalidatePartitionStatus, invalidateSchemaStatus))
{
LOG.error(
@@ -389,14 +390,14 @@ public class ConfigNodeProcedureEnv {
.sendSyncRequestToDataNodeWithGivenRetry(
dataNodeLocation.getInternalEndPoint(),
NodeStatus.Removing.getStatus(),
- CnToDnRequestType.SET_SYSTEM_STATUS,
+ CnToDnSyncRequestType.SET_SYSTEM_STATUS,
1);
} else {
SyncDataNodeClientPool.getInstance()
.sendSyncRequestToDataNodeWithRetry(
dataNodeLocation.getInternalEndPoint(),
NodeStatus.Removing.getStatus(),
- CnToDnRequestType.SET_SYSTEM_STATUS);
+ CnToDnSyncRequestType.SET_SYSTEM_STATUS);
}
long currentTime = System.nanoTime();
@@ -472,7 +473,7 @@ public class ConfigNodeProcedureEnv {
private DataNodeAsyncRequestContext<TCreateSchemaRegionReq, TSStatus>
getCreateSchemaRegionClientHandler(CreateRegionGroupsPlan
createRegionGroupsPlan) {
DataNodeAsyncRequestContext<TCreateSchemaRegionReq, TSStatus>
clientHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnRequestType.CREATE_SCHEMA_REGION);
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.CREATE_SCHEMA_REGION);
int requestId = 0;
for (Map.Entry<String, List<TRegionReplicaSet>> sgRegionsEntry :
@@ -495,7 +496,7 @@ public class ConfigNodeProcedureEnv {
private DataNodeAsyncRequestContext<TCreateDataRegionReq, TSStatus>
getCreateDataRegionClientHandler(CreateRegionGroupsPlan
createRegionGroupsPlan) {
DataNodeAsyncRequestContext<TCreateDataRegionReq, TSStatus> clientHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnRequestType.CREATE_DATA_REGION);
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.CREATE_DATA_REGION);
int requestId = 0;
for (Map.Entry<String, List<TRegionReplicaSet>> sgRegionsEntry :
@@ -590,7 +591,7 @@ public class ConfigNodeProcedureEnv {
DataNodeAsyncRequestContext<TCreateTriggerInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CREATE_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.CREATE_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -604,7 +605,7 @@ public class ConfigNodeProcedureEnv {
DataNodeAsyncRequestContext<TDropTriggerInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.DROP_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.DROP_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -617,7 +618,7 @@ public class ConfigNodeProcedureEnv {
DataNodeAsyncRequestContext<TActiveTriggerInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.ACTIVE_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.ACTIVE_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -630,7 +631,7 @@ public class ConfigNodeProcedureEnv {
DataNodeAsyncRequestContext<TInactiveTriggerInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.INACTIVE_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.INACTIVE_TRIGGER_INSTANCE, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -644,7 +645,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TCreatePipePluginInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CREATE_PIPE_PLUGIN, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.CREATE_PIPE_PLUGIN, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -658,7 +659,7 @@ public class ConfigNodeProcedureEnv {
DataNodeAsyncRequestContext<TDropPipePluginInstanceReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.DROP_PIPE_PLUGIN, request, dataNodeLocationMap);
+ CnToDnAsyncRequestType.DROP_PIPE_PLUGIN, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList();
}
@@ -671,7 +672,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushPipeMetaReq, TPushPipeMetaResp>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.PIPE_PUSH_ALL_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.PIPE_PUSH_ALL_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -686,7 +687,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushSinglePipeMetaReq,
TPushPipeMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.PIPE_PUSH_SINGLE_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.PIPE_PUSH_SINGLE_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -702,7 +703,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushSinglePipeMetaReq,
TPushPipeMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.PIPE_PUSH_SINGLE_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.PIPE_PUSH_SINGLE_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -719,7 +720,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushMultiPipeMetaReq,
TPushPipeMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.PIPE_PUSH_MULTI_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.PIPE_PUSH_MULTI_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -735,7 +736,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushMultiPipeMetaReq,
TPushPipeMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.PIPE_PUSH_MULTI_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.PIPE_PUSH_MULTI_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -751,7 +752,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushTopicMetaReq, TPushTopicMetaResp>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.TOPIC_PUSH_ALL_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.TOPIC_PUSH_ALL_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -766,7 +767,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushSingleTopicMetaReq,
TPushTopicMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.TOPIC_PUSH_SINGLE_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.TOPIC_PUSH_SINGLE_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList().stream()
.map(TPushTopicMetaResp::getStatus)
@@ -781,7 +782,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushSingleTopicMetaReq,
TPushTopicMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.TOPIC_PUSH_SINGLE_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.TOPIC_PUSH_SINGLE_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList().stream()
.map(TPushTopicMetaResp::getStatus)
@@ -797,7 +798,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushMultiTopicMetaReq,
TPushTopicMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.TOPIC_PUSH_MULTI_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.TOPIC_PUSH_MULTI_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -813,7 +814,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushMultiTopicMetaReq,
TPushTopicMetaResp> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.TOPIC_PUSH_MULTI_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.TOPIC_PUSH_MULTI_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -831,7 +832,7 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushConsumerGroupMetaReq,
TPushConsumerGroupMetaResp>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CONSUMER_GROUP_PUSH_ALL_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.CONSUMER_GROUP_PUSH_ALL_META, request,
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
.sendAsyncRequestToNodeWithRetryAndTimeoutInMs(
clientHandler,
@@ -848,7 +849,9 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushSingleConsumerGroupMetaReq,
TPushConsumerGroupMetaResp>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CONSUMER_GROUP_PUSH_SINGLE_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.CONSUMER_GROUP_PUSH_SINGLE_META,
+ request,
+ dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList().stream()
.map(TPushConsumerGroupMetaResp::getStatus)
@@ -864,7 +867,9 @@ public class ConfigNodeProcedureEnv {
final DataNodeAsyncRequestContext<TPushSingleConsumerGroupMetaReq,
TPushConsumerGroupMetaResp>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.CONSUMER_GROUP_PUSH_SINGLE_META, request,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.CONSUMER_GROUP_PUSH_SINGLE_META,
+ request,
+ dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
return clientHandler.getResponseList().stream()
.map(TPushConsumerGroupMetaResp::getStatus)
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/RegionMaintainHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/RegionMaintainHandler.java
index 490be8be342..827ecd8b6d6 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/RegionMaintainHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/RegionMaintainHandler.java
@@ -35,9 +35,10 @@ import org.apache.iotdb.commons.cluster.RegionStatus;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.commons.utils.CommonDateTimeUtils;
import org.apache.iotdb.commons.utils.NodeUrlUtils;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
+import org.apache.iotdb.confignode.client.sync.CnToDnSyncRequestType;
import org.apache.iotdb.confignode.client.sync.SyncDataNodeClientPool;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
@@ -145,7 +146,7 @@ public class RegionMaintainHandler {
.sendSyncRequestToDataNodeWithRetry(
node.getLocation().getInternalEndPoint(),
disableReq,
- CnToDnRequestType.DISABLE_DATA_NODE);
+ CnToDnSyncRequestType.DISABLE_DATANODE);
if (!isSucceed(status)) {
LOGGER.error(
"{}, BroadcastDisableDataNode meets error, disabledDataNode: {},
error: {}",
@@ -230,7 +231,7 @@ public class RegionMaintainHandler {
.sendSyncRequestToDataNodeWithRetry(
destDataNode.getInternalEndPoint(),
req,
- CnToDnRequestType.CREATE_NEW_REGION_PEER);
+ CnToDnSyncRequestType.CREATE_NEW_REGION_PEER);
if (isSucceed(status)) {
LOGGER.info(
@@ -276,7 +277,7 @@ public class RegionMaintainHandler {
.sendSyncRequestToDataNodeWithRetry(
coordinator.getInternalEndPoint(),
maintainPeerReq,
- CnToDnRequestType.ADD_REGION_PEER);
+ CnToDnSyncRequestType.ADD_REGION_PEER);
LOGGER.info(
"{}, Send action addRegionPeer finished, regionId: {}, rpcDataNode:
{}, destDataNode: {}, status: {}",
REGION_MIGRATE_PROCESS,
@@ -313,7 +314,7 @@ public class RegionMaintainHandler {
.sendSyncRequestToDataNodeWithRetry(
coordinator.getInternalEndPoint(),
maintainPeerReq,
- CnToDnRequestType.REMOVE_REGION_PEER);
+ CnToDnSyncRequestType.REMOVE_REGION_PEER);
LOGGER.info(
"{}, Send action removeRegionPeer finished, regionId: {}, rpcDataNode:
{}",
REGION_MIGRATE_PROCESS,
@@ -347,14 +348,14 @@ public class RegionMaintainHandler {
.sendSyncRequestToDataNodeWithGivenRetry(
originalDataNode.getInternalEndPoint(),
maintainPeerReq,
- CnToDnRequestType.DELETE_OLD_REGION_PEER,
+ CnToDnSyncRequestType.DELETE_OLD_REGION_PEER,
1)
: (TSStatus)
SyncDataNodeClientPool.getInstance()
.sendSyncRequestToDataNodeWithRetry(
originalDataNode.getInternalEndPoint(),
maintainPeerReq,
- CnToDnRequestType.DELETE_OLD_REGION_PEER);
+ CnToDnSyncRequestType.DELETE_OLD_REGION_PEER);
LOGGER.info(
"{}, Send action deleteOldRegionPeer finished, regionId: {},
dataNodeId: {}",
REGION_MIGRATE_PROCESS,
@@ -369,7 +370,7 @@ public class RegionMaintainHandler {
Map<Integer, TDataNodeLocation> dataNodeLocationMap) {
DataNodeAsyncRequestContext<TResetPeerListReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.RESET_PEER_LIST,
+ CnToDnAsyncRequestType.RESET_PEER_LIST,
new TResetPeerListReq(regionId, correctDataNodeLocations),
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
@@ -508,7 +509,10 @@ public class RegionMaintainHandler {
(TSStatus)
SyncDataNodeClientPool.getInstance()
.sendSyncRequestToDataNodeWithGivenRetry(
- dataNode.getInternalEndPoint(), dataNode,
CnToDnRequestType.STOP_DATA_NODE, 2);
+ dataNode.getInternalEndPoint(),
+ dataNode,
+ CnToDnSyncRequestType.STOP_DATA_NODE,
+ 2);
configManager.getLoadManager().removeNodeCache(dataNode.getDataNodeId());
LOGGER.info(
"{}, Stop Data Node result: {}, stoppedDataNode: {}",
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
index 865c1d55ad6..9009f3bffa1 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
@@ -30,7 +30,7 @@ import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.path.PathPatternTree;
import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
@@ -122,7 +122,7 @@ public class AlterLogicalViewProcedure
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TInvalidateMatchedSchemaCacheReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
new TInvalidateMatchedSchemaCacheReq(patternTreeBytes),
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
@@ -155,7 +155,7 @@ public class AlterLogicalViewProcedure
"Alter view",
env,
targetSchemaRegionGroup,
- CnToDnRequestType.ALTER_VIEW,
+ CnToDnAsyncRequestType.ALTER_VIEW,
(dataNodeLocation, consensusGroupIdList) -> {
TAlterViewReq req = new
TAlterViewReq().setIsGeneratedByPipe(isGeneratedByPipe);
req.setSchemaRegionIdList(consensusGroupIdList);
@@ -323,7 +323,7 @@ public class AlterLogicalViewProcedure
String taskName,
ConfigNodeProcedureEnv env,
Map<TConsensusGroupId, TRegionReplicaSet> targetSchemaRegionGroup,
- CnToDnRequestType dataNodeRequestType,
+ CnToDnAsyncRequestType dataNodeRequestType,
BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
super(env, targetSchemaRegionGroup, false, dataNodeRequestType,
dataNodeRequestGenerator);
this.taskName = taskName;
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DataNodeRegionTaskExecutor.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DataNodeRegionTaskExecutor.java
index af15174daf7..64a68a59c49 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DataNodeRegionTaskExecutor.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DataNodeRegionTaskExecutor.java
@@ -22,7 +22,7 @@ package org.apache.iotdb.confignode.procedure.impl.schema;
import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.manager.ConfigManager;
@@ -45,7 +45,7 @@ public abstract class DataNodeRegionTaskExecutor<Q, R> {
protected final Map<TConsensusGroupId, TRegionReplicaSet> targetRegionGroup;
protected final boolean executeOnAllReplicaset;
- protected final CnToDnRequestType dataNodeRequestType;
+ protected final CnToDnAsyncRequestType dataNodeRequestType;
protected final BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator;
@@ -55,7 +55,7 @@ public abstract class DataNodeRegionTaskExecutor<Q, R> {
ConfigManager configManager,
Map<TConsensusGroupId, TRegionReplicaSet> targetRegionGroup,
boolean executeOnAllReplicaset,
- CnToDnRequestType dataNodeRequestType,
+ CnToDnAsyncRequestType dataNodeRequestType,
BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
this.configManager = configManager;
this.targetRegionGroup = targetRegionGroup;
@@ -68,7 +68,7 @@ public abstract class DataNodeRegionTaskExecutor<Q, R> {
ConfigNodeProcedureEnv env,
Map<TConsensusGroupId, TRegionReplicaSet> targetRegionGroup,
boolean executeOnAllReplicaset,
- CnToDnRequestType dataNodeRequestType,
+ CnToDnAsyncRequestType dataNodeRequestType,
BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
this.configManager = env.getConfigManager();
this.targetRegionGroup = targetRegionGroup;
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeactivateTemplateProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeactivateTemplateProcedure.java
index dd3cb6b0073..24971e1ea36 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeactivateTemplateProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeactivateTemplateProcedure.java
@@ -28,7 +28,7 @@ import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.path.PathPatternTree;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.write.pipe.payload.PipeDeactivateTemplatePlan;
@@ -152,7 +152,7 @@ public class DeactivateTemplateProcedure
"construct schema black list",
env,
targetSchemaRegionGroup,
- CnToDnRequestType.CONSTRUCT_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
+
CnToDnAsyncRequestType.CONSTRUCT_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
((dataNodeLocation, consensusGroupIdList) ->
new TConstructSchemaBlackListWithTemplateReq(
consensusGroupIdList, dataNodeRequest))) {
@@ -200,7 +200,7 @@ public class DeactivateTemplateProcedure
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TInvalidateMatchedSchemaCacheReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
new TInvalidateMatchedSchemaCacheReq(timeSeriesPatternTreeBytes),
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance()
@@ -233,7 +233,7 @@ public class DeactivateTemplateProcedure
env,
relatedDataRegionGroup,
true,
- CnToDnRequestType.DELETE_DATA_FOR_DELETE_SCHEMA,
+ CnToDnAsyncRequestType.DELETE_DATA_FOR_DELETE_SCHEMA,
((dataNodeLocation, consensusGroupIdList) ->
new TDeleteDataForDeleteSchemaReq(
new ArrayList<>(consensusGroupIdList),
timeSeriesPatternTreeBytes)));
@@ -248,7 +248,7 @@ public class DeactivateTemplateProcedure
"deactivate template schema",
env,
env.getConfigManager().getRelatedSchemaRegionGroup(timeSeriesPatternTree),
- CnToDnRequestType.DEACTIVATE_TEMPLATE,
+ CnToDnAsyncRequestType.DEACTIVATE_TEMPLATE,
((dataNodeLocation, consensusGroupIdList) ->
new TDeactivateTemplateReq(consensusGroupIdList,
dataNodeRequest)
.setIsGeneratedByPipe(isGeneratedByPipe)));
@@ -286,7 +286,7 @@ public class DeactivateTemplateProcedure
"roll back schema black list",
env,
env.getConfigManager().getRelatedSchemaRegionGroup(timeSeriesPatternTree),
- CnToDnRequestType.ROLLBACK_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
+
CnToDnAsyncRequestType.ROLLBACK_SCHEMA_BLACK_LIST_WITH_TEMPLATE,
((dataNodeLocation, consensusGroupIdList) ->
new TRollbackSchemaBlackListWithTemplateReq(
consensusGroupIdList, dataNodeRequest)));
@@ -439,7 +439,7 @@ public class DeactivateTemplateProcedure
String taskName,
ConfigNodeProcedureEnv env,
Map<TConsensusGroupId, TRegionReplicaSet> targetSchemaRegionGroup,
- CnToDnRequestType dataNodeRequestType,
+ CnToDnAsyncRequestType dataNodeRequestType,
BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
super(env, targetSchemaRegionGroup, false, dataNodeRequestType,
dataNodeRequestGenerator);
this.taskName = taskName;
@@ -450,7 +450,7 @@ public class DeactivateTemplateProcedure
ConfigNodeProcedureEnv env,
Map<TConsensusGroupId, TRegionReplicaSet> targetDataRegionGroup,
boolean executeOnAllReplicaset,
- CnToDnRequestType dataNodeRequestType,
+ CnToDnAsyncRequestType dataNodeRequestType,
BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
super(
env,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
index d47e575a3ad..dd52415d347 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
@@ -27,7 +27,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.exception.runtime.ThriftSerDeException;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.commons.utils.ThriftConfigNodeSerDeUtils;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.write.database.PreDeleteDatabasePlan;
@@ -156,7 +156,7 @@ public class DeleteDatabaseProcedure
// try sync delete schemaengine region
DataNodeAsyncRequestContext<TConsensusGroupId, TSStatus>
asyncClientHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnRequestType.DELETE_REGION);
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION);
Map<Integer, RegionDeleteTask> schemaRegionDeleteTaskMap = new
HashMap<>();
int requestIndex = 0;
for (TRegionReplicaSet schemaRegionReplicaSet :
schemaRegionReplicaSets) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
index 0a83770689c..ed680d9ca38 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
@@ -26,7 +26,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathPatternTree;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.write.pipe.payload.PipeDeleteLogicalViewPlan;
@@ -142,7 +142,7 @@ public class DeleteLogicalViewProcedure
"construct view schema engine black list",
env,
targetSchemaRegionGroup,
- CnToDnRequestType.CONSTRUCT_VIEW_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.CONSTRUCT_VIEW_SCHEMA_BLACK_LIST,
((dataNodeLocation, consensusGroupIdList) ->
new TConstructViewSchemaBlackListReq(consensusGroupIdList,
patternTreeBytes))) {
@Override
@@ -186,7 +186,7 @@ public class DeleteLogicalViewProcedure
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TInvalidateMatchedSchemaCacheReq, TSStatus>
clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
new TInvalidateMatchedSchemaCacheReq(patternTreeBytes),
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
@@ -211,7 +211,7 @@ public class DeleteLogicalViewProcedure
"delete view in schema engine",
env,
env.getConfigManager().getRelatedSchemaRegionGroup(patternTree),
- CnToDnRequestType.DELETE_VIEW,
+ CnToDnAsyncRequestType.DELETE_VIEW,
((dataNodeLocation, consensusGroupIdList) ->
new TDeleteViewSchemaReq(consensusGroupIdList,
patternTreeBytes)
.setIsGeneratedByPipe(isGeneratedByPipe)));
@@ -247,7 +247,7 @@ public class DeleteLogicalViewProcedure
"roll back view schema engine black list",
env,
env.getConfigManager().getRelatedSchemaRegionGroup(patternTree),
- CnToDnRequestType.ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
(dataNodeLocation, consensusGroupIdList) ->
new TRollbackViewSchemaBlackListReq(consensusGroupIdList,
patternTreeBytes));
rollbackStateTask.execute();
@@ -343,7 +343,7 @@ public class DeleteLogicalViewProcedure
String taskName,
ConfigNodeProcedureEnv env,
Map<TConsensusGroupId, TRegionReplicaSet> targetSchemaRegionGroup,
- CnToDnRequestType dataNodeRequestType,
+ CnToDnAsyncRequestType dataNodeRequestType,
BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
super(env, targetSchemaRegionGroup, false, dataNodeRequestType,
dataNodeRequestGenerator);
this.taskName = taskName;
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteTimeSeriesProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteTimeSeriesProcedure.java
index cd9133c0cf5..bfb1dc7c19b 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteTimeSeriesProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteTimeSeriesProcedure.java
@@ -26,7 +26,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathPatternTree;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.write.pipe.payload.PipeDeleteTimeSeriesPlan;
@@ -154,7 +154,7 @@ public class DeleteTimeSeriesProcedure
"construct schema engine black list",
env,
targetSchemaRegionGroup,
- CnToDnRequestType.CONSTRUCT_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.CONSTRUCT_SCHEMA_BLACK_LIST,
((dataNodeLocation, consensusGroupIdList) ->
new TConstructSchemaBlackListReq(consensusGroupIdList,
patternTreeBytes))) {
@Override
@@ -200,7 +200,7 @@ public class DeleteTimeSeriesProcedure
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
final DataNodeAsyncRequestContext<TInvalidateMatchedSchemaCacheReq,
TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
+ CnToDnAsyncRequestType.INVALIDATE_MATCHED_SCHEMA_CACHE,
new TInvalidateMatchedSchemaCacheReq(patternTreeBytes),
dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
@@ -250,7 +250,7 @@ public class DeleteTimeSeriesProcedure
env,
relatedDataRegionGroup,
true,
- CnToDnRequestType.DELETE_DATA_FOR_DELETE_SCHEMA,
+ CnToDnAsyncRequestType.DELETE_DATA_FOR_DELETE_SCHEMA,
((dataNodeLocation, consensusGroupIdList) ->
new TDeleteDataForDeleteSchemaReq(
new ArrayList<>(consensusGroupIdList),
@@ -265,7 +265,7 @@ public class DeleteTimeSeriesProcedure
"delete time series in schema engine",
env,
env.getConfigManager().getRelatedSchemaRegionGroup(patternTree),
- CnToDnRequestType.DELETE_TIMESERIES,
+ CnToDnAsyncRequestType.DELETE_TIMESERIES,
((dataNodeLocation, consensusGroupIdList) ->
new TDeleteTimeSeriesReq(consensusGroupIdList,
patternTreeBytes)
.setIsGeneratedByPipe(isGeneratedByPipe)));
@@ -302,7 +302,7 @@ public class DeleteTimeSeriesProcedure
"roll back schema engine black list",
env,
env.getConfigManager().getRelatedSchemaRegionGroup(patternTree),
- CnToDnRequestType.ROLLBACK_SCHEMA_BLACK_LIST,
+ CnToDnAsyncRequestType.ROLLBACK_SCHEMA_BLACK_LIST,
(dataNodeLocation, consensusGroupIdList) ->
new TRollbackSchemaBlackListReq(consensusGroupIdList,
patternTreeBytes));
rollbackStateTask.execute();
@@ -403,7 +403,7 @@ public class DeleteTimeSeriesProcedure
final String taskName,
final ConfigNodeProcedureEnv env,
final Map<TConsensusGroupId, TRegionReplicaSet>
targetSchemaRegionGroup,
- final CnToDnRequestType dataNodeRequestType,
+ final CnToDnAsyncRequestType dataNodeRequestType,
final BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
super(env, targetSchemaRegionGroup, false, dataNodeRequestType,
dataNodeRequestGenerator);
this.taskName = taskName;
@@ -414,7 +414,7 @@ public class DeleteTimeSeriesProcedure
final ConfigNodeProcedureEnv env,
final Map<TConsensusGroupId, TRegionReplicaSet> targetDataRegionGroup,
final boolean executeOnAllReplicaset,
- final CnToDnRequestType dataNodeRequestType,
+ final CnToDnAsyncRequestType dataNodeRequestType,
final BiFunction<TDataNodeLocation, List<TConsensusGroupId>, Q>
dataNodeRequestGenerator) {
super(
env,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SchemaUtils.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SchemaUtils.java
index 1df6ff508fd..25cd4704dc6 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SchemaUtils.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SchemaUtils.java
@@ -26,7 +26,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathPatternTree;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.db.exception.metadata.PathNotExistException;
import org.apache.iotdb.db.schemaengine.template.Template;
@@ -78,7 +78,7 @@ public class SchemaUtils {
configManager,
relatedSchemaRegionGroup,
false,
- CnToDnRequestType.COUNT_PATHS_USING_TEMPLATE,
+ CnToDnAsyncRequestType.COUNT_PATHS_USING_TEMPLATE,
((dataNodeLocation, consensusGroupIdList) ->
new TCountPathsUsingTemplateReq(
template.getId(), patternTreeBytes,
consensusGroupIdList))) {
@@ -156,7 +156,7 @@ public class SchemaUtils {
configManager,
relatedSchemaRegionGroup,
false,
- CnToDnRequestType.CHECK_SCHEMA_REGION_USING_TEMPLATE,
+ CnToDnAsyncRequestType.CHECK_SCHEMA_REGION_USING_TEMPLATE,
((dataNodeLocation, consensusGroupIdList) ->
new
TCheckSchemaRegionUsingTemplateReq(consensusGroupIdList))) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTTLProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTTLProcedure.java
index ec9b003d7ee..de81a0681a1 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTTLProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTTLProcedure.java
@@ -24,7 +24,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.common.rpc.thrift.TSetTTLReq;
import org.apache.iotdb.commons.exception.IoTDBException;
import org.apache.iotdb.commons.exception.MetadataException;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan;
@@ -110,7 +110,7 @@ public class SetTTLProcedure extends
StateMachineProcedure<ConfigNodeProcedureEn
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TSetTTLReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.SET_TTL,
+ CnToDnAsyncRequestType.SET_TTL,
new TSetTTLReq(
Collections.singletonList(String.join(".",
plan.getPathPattern())),
plan.getTTL(),
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTemplateProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTemplateProcedure.java
index 2bd4068ee89..8a9f1403d7d 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTemplateProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/SetTemplateProcedure.java
@@ -28,7 +28,7 @@ import org.apache.iotdb.commons.exception.IoTDBException;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathPatternTree;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import
org.apache.iotdb.confignode.consensus.request.read.template.CheckTemplateSettablePlan;
@@ -212,7 +212,7 @@ public class SetTemplateProcedure
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TUpdateTemplateReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TEMPLATE, req, dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE, req, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
Map<Integer, TSStatus> statusMap = clientHandler.getResponseMap();
for (Map.Entry<Integer, TSStatus> entry : statusMap.entrySet()) {
@@ -282,7 +282,7 @@ public class SetTemplateProcedure
env,
relatedSchemaRegionGroup,
false,
- CnToDnRequestType.CHECK_TIMESERIES_EXISTENCE,
+ CnToDnAsyncRequestType.CHECK_TIMESERIES_EXISTENCE,
((dataNodeLocation, consensusGroupIdList) ->
new TCheckTimeSeriesExistenceReq(patternTreeBytes,
consensusGroupIdList))) {
@@ -386,7 +386,7 @@ public class SetTemplateProcedure
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
DataNodeAsyncRequestContext<TUpdateTemplateReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TEMPLATE, req, dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE, req, dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
Map<Integer, TSStatus> statusMap = clientHandler.getResponseMap();
for (Map.Entry<Integer, TSStatus> entry : statusMap.entrySet()) {
@@ -491,7 +491,9 @@ public class SetTemplateProcedure
DataNodeAsyncRequestContext<TUpdateTemplateReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TEMPLATE, invalidateTemplateSetInfoReq,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE,
+ invalidateTemplateSetInfoReq,
+ dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
Map<Integer, TSStatus> statusMap = clientHandler.getResponseMap();
for (Map.Entry<Integer, TSStatus> entry : statusMap.entrySet()) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/UnsetTemplateProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/UnsetTemplateProcedure.java
index ee70f3c7a7d..d85d1bfb301 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/UnsetTemplateProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/UnsetTemplateProcedure.java
@@ -26,7 +26,7 @@ import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.path.PathPatternTree;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
import
org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager;
import
org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
@@ -157,7 +157,9 @@ public class UnsetTemplateProcedure
invalidateTemplateSetInfoReq.setTemplateInfo(getInvalidateTemplateSetInfo());
DataNodeAsyncRequestContext<TUpdateTemplateReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TEMPLATE, invalidateTemplateSetInfoReq,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE,
+ invalidateTemplateSetInfoReq,
+ dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
Map<Integer, TSStatus> statusMap = clientHandler.getResponseMap();
for (TSStatus status : statusMap.values()) {
@@ -252,7 +254,9 @@ public class UnsetTemplateProcedure
rollbackTemplateSetInfoReq.setTemplateInfo(getAddTemplateSetInfo());
DataNodeAsyncRequestContext<TUpdateTemplateReq, TSStatus> clientHandler =
new DataNodeAsyncRequestContext<>(
- CnToDnRequestType.UPDATE_TEMPLATE, rollbackTemplateSetInfoReq,
dataNodeLocationMap);
+ CnToDnAsyncRequestType.UPDATE_TEMPLATE,
+ rollbackTemplateSetInfoReq,
+ dataNodeLocationMap);
CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequestWithRetry(clientHandler);
Map<Integer, TSStatus> statusMap = clientHandler.getResponseMap();
for (TSStatus status : statusMap.values()) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/sync/AuthOperationProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/sync/AuthOperationProcedure.java
index f8e0c9d2233..20cc759a74c 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/sync/AuthOperationProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/sync/AuthOperationProcedure.java
@@ -25,7 +25,7 @@ import org.apache.iotdb.commons.conf.CommonConfig;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.exception.IoTDBException;
import org.apache.iotdb.commons.utils.ThriftCommonsSerDeUtils;
-import org.apache.iotdb.confignode.client.CnToDnRequestType;
+import org.apache.iotdb.confignode.client.sync.CnToDnSyncRequestType;
import org.apache.iotdb.confignode.client.sync.SyncDataNodeClientPool;
import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan;
import org.apache.iotdb.confignode.consensus.request.write.auth.AuthorPlan;
@@ -112,7 +112,7 @@ public class AuthOperationProcedure extends
AbstractNodeProcedure<AuthOperationP
.sendSyncRequestToDataNodeWithRetry(
pair.getLeft().getLocation().getInternalEndPoint(),
req,
- CnToDnRequestType.INVALIDATE_PERMISSION_CACHE);
+ CnToDnSyncRequestType.INVALIDATE_PERMISSION_CACHE);
if (status.getCode() ==
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
it.remove();
}
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/request/AsyncRequestManager.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/request/AsyncRequestManager.java
index e0a2ea0b635..9ca809a4d8c 100644
---
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/request/AsyncRequestManager.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/request/AsyncRequestManager.java
@@ -58,12 +58,18 @@ public abstract class AsyncRequestManager<RequestType,
NodeLocation, Client> {
actionMapBuilder = ImmutableMap.builder();
initActionMapBuilder();
this.actionMap = this.actionMapBuilder.build();
+ checkActionMapCompleteness();
}
protected abstract void initClientManager();
protected abstract void initActionMapBuilder();
+ protected void checkActionMapCompleteness() {
+ // No check by default
+ }
+ ;
+
/**
* Send asynchronous requests to the specified Nodes with default retry num
*
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/exception/IoTDBRuntimeException.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/exception/IoTDBRuntimeException.java
new file mode 100644
index 00000000000..3d58f7d606b
--- /dev/null
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/exception/IoTDBRuntimeException.java
@@ -0,0 +1,66 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ *
+ */
+
+package org.apache.iotdb.commons.exception;
+
+public class IoTDBRuntimeException extends RuntimeException {
+ protected int errorCode;
+
+ /**
+ * This kind of exception is caused by users' wrong sql, and there is no
need for server to print
+ * the full stack of the exception
+ */
+ protected boolean isUserException = false;
+
+ public IoTDBRuntimeException(String message, int errorCode) {
+ super(message);
+ this.errorCode = errorCode;
+ }
+
+ public IoTDBRuntimeException(String message, int errorCode, boolean
isUserException) {
+ super(message);
+ this.errorCode = errorCode;
+ this.isUserException = isUserException;
+ }
+
+ public IoTDBRuntimeException(String message, Throwable cause, int errorCode)
{
+ super(message, cause);
+ this.errorCode = errorCode;
+ }
+
+ public IoTDBRuntimeException(Throwable cause, int errorCode) {
+ super(cause);
+ this.errorCode = errorCode;
+ }
+
+ public IoTDBRuntimeException(Throwable cause, int errorCode, boolean
isUserException) {
+ super(cause);
+ this.errorCode = errorCode;
+ this.isUserException = isUserException;
+ }
+
+ public boolean isUserException() {
+ return isUserException;
+ }
+
+ public int getErrorCode() {
+ return errorCode;
+ }
+}
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/exception/UncheckedStartupException.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/exception/UncheckedStartupException.java
new file mode 100644
index 00000000000..46649b9385b
--- /dev/null
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/exception/UncheckedStartupException.java
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.commons.exception;
+
+import org.apache.iotdb.rpc.TSStatusCode;
+
+public class UncheckedStartupException extends IoTDBRuntimeException {
+
+ private static final long serialVersionUID = -8591716406230730147L;
+
+ public UncheckedStartupException(String name, String message) {
+ super(
+ String.format("Failed to start [%s], because [%s]", name, message),
+ TSStatusCode.START_UP_ERROR.getStatusCode());
+ }
+
+ public UncheckedStartupException(Throwable cause) {
+ super(cause.getMessage(), TSStatusCode.START_UP_ERROR.getStatusCode());
+ }
+
+ public UncheckedStartupException(String message) {
+ super(message, TSStatusCode.START_UP_ERROR.getStatusCode());
+ }
+}