This is an automated email from the ASF dual-hosted git repository.
rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 478e4d17952 [IOTDB-5839] Pipe task management (CN -> DN): squash all
operation rpcs into one (#9750)
478e4d17952 is described below
commit 478e4d1795287823bffc16ac71acadbd88723128
Author: Caideyipi <[email protected]>
AuthorDate: Sun May 7 12:11:42 2023 +0800
[IOTDB-5839] Pipe task management (CN -> DN): squash all operation rpcs
into one (#9750)
Co-authored-by: Steve Yurong Su <[email protected]>
---
.../confignode/client/DataNodeRequestType.java | 7 +--
.../client/async/AsyncDataNodeClientPool.java | 15 ++----
.../client/async/handlers/AsyncClientHandler.java | 1 +
.../confignode/persistence/pipe/PipeTaskInfo.java | 4 ++
.../procedure/env/ConfigNodeProcedureEnv.java | 23 ++------
.../pipe/task/AbstractOperatePipeProcedureV2.java | 43 ++++++++++++++-
.../impl/pipe/task/CreatePipeProcedureV2.java | 28 ++--------
.../impl/pipe/task/DropPipeProcedureV2.java | 16 ++----
.../impl/pipe/task/StartPipeProcedureV2.java | 28 ++--------
.../impl/pipe/task/StopPipeProcedureV2.java | 28 ++--------
.../iotdb/confignode/persistence/PipeInfoTest.java | 1 +
.../iotdb/commons/pipe/task/meta/PipeMeta.java | 10 +++-
.../commons/pipe/task/meta/PipeMetaKeeper.java | 4 ++
.../commons/pipe/task/meta/PipeRuntimeMeta.java | 20 +------
.../commons/pipe/task/meta/PipeStaticMeta.java | 63 +++++++++++-----------
.../iotdb/commons/pipe/task/meta/PipeTaskMeta.java | 51 ++++++++++++++++--
.../iotdb/pipe/api/customizer/PipeParameters.java | 22 ++++++++
.../exception/PipeRuntimeCriticalException.java | 40 ++++++++++++++
.../pipe/api/exception/PipeRuntimeException.java | 40 ++++++++++++++
.../exception/PipeRuntimeNonCriticalException.java | 40 ++++++++++++++
.../impl/DataNodeInternalRPCServiceImpl.java | 31 +++--------
thrift/src/main/thrift/datanode.thrift | 19 ++-----
22 files changed, 318 insertions(+), 216 deletions(-)
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
index 56ec55d53f1..d36dce571f6 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
@@ -67,12 +67,7 @@ public enum DataNodeRequestType {
DROP_PIPE_PLUGIN,
/** Pipe Task */
- CREATE_PIPE,
- /**
- * DROP_PIPE, START_PIPE, STOP_PIPE Merge them into OPERATE_PIPE since these
requests only require
- * pipe name
- */
- OPERATE_PIPE,
+ PUSH_PIPE_META,
/** CQ */
EXECUTE_CQ,
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
index ace4d0268b8..332b7e2317b 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
@@ -44,7 +44,6 @@ import
org.apache.iotdb.mpp.rpc.thrift.TConstructSchemaBlackListWithTemplateReq;
import org.apache.iotdb.mpp.rpc.thrift.TCountPathsUsingTemplateReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateDataRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateFunctionInstanceReq;
-import org.apache.iotdb.mpp.rpc.thrift.TCreatePipeOnDataNodeReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreatePipePluginInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateSchemaRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateTriggerInstanceReq;
@@ -57,7 +56,7 @@ import
org.apache.iotdb.mpp.rpc.thrift.TDropTriggerInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TFetchSchemaBlackListReq;
import org.apache.iotdb.mpp.rpc.thrift.TInactiveTriggerInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TInvalidateMatchedSchemaCacheReq;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
+import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaReq;
import org.apache.iotdb.mpp.rpc.thrift.TRegionLeaderChangeReq;
import org.apache.iotdb.mpp.rpc.thrift.TRegionRouteReq;
import org.apache.iotdb.mpp.rpc.thrift.TRollbackSchemaBlackListReq;
@@ -221,9 +220,9 @@ public class AsyncDataNodeClientPool {
(AsyncTSStatusRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
- case CREATE_PIPE:
- client.createPipeOnDataNode(
- (TCreatePipeOnDataNodeReq) clientHandler.getRequest(requestId),
+ case PUSH_PIPE_META:
+ client.pushPipeMeta(
+ (TPushPipeMetaReq) clientHandler.getRequest(requestId),
(AsyncTSStatusRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
@@ -309,12 +308,6 @@ public class AsyncDataNodeClientPool {
(DeleteSchemaRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
- case OPERATE_PIPE:
- client.operatePipeOnDataNode(
- (TOperatePipeOnDataNodeReq) clientHandler.getRequest(requestId),
- (AsyncTSStatusRPCHandler)
- clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
- break;
case CONSTRUCT_SCHEMA_BLACK_LIST_WITH_TEMPLATE:
client.constructSchemaBlackListWithTemplate(
(TConstructSchemaBlackListWithTemplateReq)
clientHandler.getRequest(requestId),
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
index f54c0e3977b..a5904b5f8a8 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
@@ -215,6 +215,7 @@ public class AsyncClientHandler<Q, R> {
case UPDATE_TEMPLATE:
case CHANGE_REGION_LEADER:
case KILL_QUERY_INSTANCE:
+ case PUSH_PIPE_META:
default:
return new AsyncTSStatusRPCHandler(
requestType,
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
index 9ea25320ff3..c3975a817d7 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
@@ -144,6 +144,10 @@ public class PipeTaskInfo implements SnapshotProcessor {
return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
}
+ public Iterable<PipeMeta> getPipeMetaList() {
+ return pipeMetaKeeper.getPipeMetaList();
+ }
+
/////////////////////////////// Snapshot ///////////////////////////////
@Override
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
index c7a500f5dc9..62dd2b538e1 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
@@ -30,7 +30,6 @@ import org.apache.iotdb.commons.cluster.NodeStatus;
import org.apache.iotdb.commons.cluster.NodeType;
import org.apache.iotdb.commons.cluster.RegionStatus;
import org.apache.iotdb.commons.pipe.plugin.meta.PipePluginMeta;
-import org.apache.iotdb.commons.pipe.task.meta.PipeMeta;
import org.apache.iotdb.commons.trigger.TriggerInformation;
import org.apache.iotdb.confignode.client.ConfigNodeRequestType;
import org.apache.iotdb.confignode.client.DataNodeRequestType;
@@ -60,7 +59,6 @@ import
org.apache.iotdb.confignode.procedure.scheduler.ProcedureScheduler;
import org.apache.iotdb.confignode.rpc.thrift.TAddConsensusGroupReq;
import org.apache.iotdb.mpp.rpc.thrift.TActiveTriggerInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateDataRegionReq;
-import org.apache.iotdb.mpp.rpc.thrift.TCreatePipeOnDataNodeReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreatePipePluginInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateSchemaRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateTriggerInstanceReq;
@@ -68,7 +66,7 @@ import
org.apache.iotdb.mpp.rpc.thrift.TDropPipePluginInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TDropTriggerInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TInactiveTriggerInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TInvalidateCacheReq;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
+import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaReq;
import org.apache.iotdb.mpp.rpc.thrift.TUpdateConfigNodeGroupReq;
import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.Binary;
@@ -638,24 +636,13 @@ public class ConfigNodeProcedureEnv {
return clientHandler.getResponseList();
}
- public List<TSStatus> createPipeOnDataNodes(PipeMeta pipeMeta) throws
IOException {
+ public List<TSStatus> pushPipeMetaToDataNodes(List<ByteBuffer>
pipeMetaBinaryList) {
final Map<Integer, TDataNodeLocation> dataNodeLocationMap =
configManager.getNodeManager().getRegisteredDataNodeLocations();
+ final TPushPipeMetaReq request = new
TPushPipeMetaReq().setPipeMetas(pipeMetaBinaryList);
- TCreatePipeOnDataNodeReq request =
- new TCreatePipeOnDataNodeReq().setPipeMeta(pipeMeta.serialize());
- final AsyncClientHandler<TCreatePipeOnDataNodeReq, TSStatus> clientHandler
=
- new AsyncClientHandler<>(DataNodeRequestType.CREATE_PIPE, request,
dataNodeLocationMap);
-
AsyncDataNodeClientPool.getInstance().sendAsyncRequestToDataNodeWithRetry(clientHandler);
- return clientHandler.getResponseList();
- }
-
- public List<TSStatus> operatePipeOnDataNodes(TOperatePipeOnDataNodeReq
request) {
- final Map<Integer, TDataNodeLocation> dataNodeLocationMap =
- configManager.getNodeManager().getRegisteredDataNodeLocations();
-
- final AsyncClientHandler<TOperatePipeOnDataNodeReq, TSStatus>
clientHandler =
- new AsyncClientHandler<>(DataNodeRequestType.OPERATE_PIPE, request,
dataNodeLocationMap);
+ final AsyncClientHandler<TPushPipeMetaReq, TSStatus> clientHandler =
+ new AsyncClientHandler<>(DataNodeRequestType.PUSH_PIPE_META, request,
dataNodeLocationMap);
AsyncDataNodeClientPool.getInstance().sendAsyncRequestToDataNodeWithRetry(clientHandler);
return clientHandler.getResponseList();
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/AbstractOperatePipeProcedureV2.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/AbstractOperatePipeProcedureV2.java
index b840d94591c..0d9f98101f8 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/AbstractOperatePipeProcedureV2.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/AbstractOperatePipeProcedureV2.java
@@ -20,6 +20,7 @@ package org.apache.iotdb.confignode.procedure.impl.pipe.task;
import org.apache.iotdb.commons.exception.sync.PipeException;
import org.apache.iotdb.commons.exception.sync.PipeSinkException;
+import org.apache.iotdb.commons.pipe.task.meta.PipeMeta;
import org.apache.iotdb.confignode.persistence.pipe.PipeTaskOperation;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.exception.ProcedureException;
@@ -27,11 +28,17 @@ import
org.apache.iotdb.confignode.procedure.exception.ProcedureSuspendedExcepti
import org.apache.iotdb.confignode.procedure.exception.ProcedureYieldException;
import org.apache.iotdb.confignode.procedure.impl.node.AbstractNodeProcedure;
import
org.apache.iotdb.confignode.procedure.state.pipe.task.OperatePipeTaskState;
+import org.apache.iotdb.pipe.api.exception.PipeManagementException;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.List;
/**
* This procedure manage 4 kinds of PIPE operations: CREATE, START, STOP and
DROP.
@@ -46,6 +53,9 @@ abstract class AbstractOperatePipeProcedureV2 extends
AbstractNodeProcedure<Oper
private static final int RETRY_THRESHOLD = 3;
+ // only used in rollback to reduce the number of network calls
+ protected boolean isRollbackFromOperateOnDataNodesSuccessful = false;
+
abstract PipeTaskOperation getOperation();
/**
@@ -126,10 +136,21 @@ abstract class AbstractOperatePipeProcedureV2 extends
AbstractNodeProcedure<Oper
rollbackFromCalculateInfoForTask(env);
break;
case WRITE_CONFIG_NODE_CONSENSUS:
- rollbackFromWriteConfigNodeConsensus(env);
+ // rollbackFromWriteConfigNodeConsensus can be called before
rollbackFromOperateOnDataNodes
+ // so we need to check if rollbackFromOperateOnDataNodes is successful
executed
+ // if yes, we don't need to call rollbackFromWriteConfigNodeConsensus
again
+ if (!isRollbackFromOperateOnDataNodesSuccessful) {
+ rollbackFromWriteConfigNodeConsensus(env);
+ }
break;
case OPERATE_ON_DATA_NODES:
+ // we have to make sure that rollbackFromOperateOnDataNodes is
executed before
+ // rollbackFromWriteConfigNodeConsensus, because
rollbackFromOperateOnDataNodes is
+ // executed based on the consensus of config nodes that is written by
+ // rollbackFromWriteConfigNodeConsensus
+ rollbackFromWriteConfigNodeConsensus(env);
rollbackFromOperateOnDataNodes(env);
+ isRollbackFromOperateOnDataNodesSuccessful = true;
break;
default:
LOGGER.error("Unsupported roll back STATE [{}]", state);
@@ -142,7 +163,8 @@ abstract class AbstractOperatePipeProcedureV2 extends
AbstractNodeProcedure<Oper
protected abstract void
rollbackFromWriteConfigNodeConsensus(ConfigNodeProcedureEnv env);
- protected abstract void
rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env);
+ protected abstract void
rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+ throws IOException;
@Override
protected OperatePipeTaskState getState(int stateId) {
@@ -158,4 +180,21 @@ abstract class AbstractOperatePipeProcedureV2 extends
AbstractNodeProcedure<Oper
protected OperatePipeTaskState getInitialState() {
return OperatePipeTaskState.VALIDATE_TASK;
}
+
+ protected void pushPipeMetaToDataNodes(ConfigNodeProcedureEnv env) throws
IOException {
+ final List<ByteBuffer> pipeMetaBinaryList = new ArrayList<>();
+ for (PipeMeta pipeMeta :
+ env.getConfigManager()
+ .getPipeManager()
+ .getPipeTaskCoordinator()
+ .getPipeTaskInfo()
+ .getPipeMetaList()) {
+ pipeMetaBinaryList.add(pipeMeta.serialize());
+ }
+
+ if
(RpcUtils.squashResponseStatusList(env.pushPipeMetaToDataNodes(pipeMetaBinaryList)).getCode()
+ != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ throw new PipeManagementException("Failed to push pipe meta list to data
nodes");
+ }
+ }
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/CreatePipeProcedureV2.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/CreatePipeProcedureV2.java
index 13f16f9ae89..995b2a977dd 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/CreatePipeProcedureV2.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/CreatePipeProcedureV2.java
@@ -20,7 +20,6 @@
package org.apache.iotdb.confignode.procedure.impl.pipe.task;
import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
-import org.apache.iotdb.commons.pipe.task.meta.PipeMeta;
import org.apache.iotdb.commons.pipe.task.meta.PipeRuntimeMeta;
import org.apache.iotdb.commons.pipe.task.meta.PipeStaticMeta;
import org.apache.iotdb.commons.pipe.task.meta.PipeTaskMeta;
@@ -32,11 +31,8 @@ import
org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
import org.apache.iotdb.confignode.rpc.thrift.TCreatePipeReq;
import org.apache.iotdb.consensus.common.response.ConsensusWriteResponse;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.iotdb.pipe.api.exception.PipeManagementException;
-import org.apache.iotdb.rpc.RpcUtils;
-import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import org.slf4j.Logger;
@@ -136,15 +132,7 @@ public class CreatePipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
"CreatePipeProcedureV2: executeFromOperateOnDataNodes({})",
createPipeRequest.getPipeName());
- if (RpcUtils.squashResponseStatusList(
- env.createPipeOnDataNodes(new PipeMeta(pipeStaticMeta,
pipeRuntimeMeta)))
- .getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format(
- "Failed to create pipe instance [%s] on data nodes",
- createPipeRequest.getPipeName()));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
@@ -178,22 +166,12 @@ public class CreatePipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
}
@Override
- protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env) {
+ protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
throws IOException {
LOGGER.info(
"CreatePipeProcedureV2: rollbackFromOperateOnDataNodes({})",
createPipeRequest.getPipeName());
- final TOperatePipeOnDataNodeReq request =
- new TOperatePipeOnDataNodeReq()
- .setPipeName(createPipeRequest.getPipeName())
- .setOperation((byte) PipeTaskOperation.DROP_PIPE.ordinal());
- if
(RpcUtils.squashResponseStatusList(env.operatePipeOnDataNodes(request)).getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format(
- "Failed to rollback from operate on data nodes for task [%s]",
- createPipeRequest.getPipeName()));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java
index 57675c7c44d..8b03d4ba486 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/DropPipeProcedureV2.java
@@ -23,11 +23,8 @@ import
org.apache.iotdb.confignode.persistence.pipe.PipeTaskOperation;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
import org.apache.iotdb.consensus.common.response.ConsensusWriteResponse;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.iotdb.pipe.api.exception.PipeManagementException;
-import org.apache.iotdb.rpc.RpcUtils;
-import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import org.slf4j.Logger;
@@ -87,18 +84,11 @@ public class DropPipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
}
@Override
- void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws
PipeManagementException {
+ void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+ throws PipeManagementException, IOException {
LOGGER.info("DropPipeProcedureV2: executeFromOperateOnDataNodes({})",
pipeName);
- final TOperatePipeOnDataNodeReq request =
- new TOperatePipeOnDataNodeReq()
- .setPipeName(pipeName)
- .setOperation((byte) PipeTaskOperation.DROP_PIPE.ordinal());
- if
(RpcUtils.squashResponseStatusList(env.operatePipeOnDataNodes(request)).getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format("Failed to drop pipe instance [%s] on data nodes",
pipeName));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StartPipeProcedureV2.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StartPipeProcedureV2.java
index 0a6be55a2c7..e1b134879b7 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StartPipeProcedureV2.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StartPipeProcedureV2.java
@@ -24,11 +24,8 @@ import
org.apache.iotdb.confignode.persistence.pipe.PipeTaskOperation;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
import org.apache.iotdb.consensus.common.response.ConsensusWriteResponse;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.iotdb.pipe.api.exception.PipeManagementException;
-import org.apache.iotdb.rpc.RpcUtils;
-import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import org.slf4j.Logger;
@@ -90,18 +87,11 @@ public class StartPipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
}
@Override
- void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws
PipeManagementException {
+ void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+ throws PipeManagementException, IOException {
LOGGER.info("StartPipeProcedureV2: executeFromOperateOnDataNodes({})",
pipeName);
- final TOperatePipeOnDataNodeReq request =
- new TOperatePipeOnDataNodeReq()
- .setPipeName(pipeName)
- .setOperation((byte) PipeTaskOperation.START_PIPE.ordinal());
- if
(RpcUtils.squashResponseStatusList(env.operatePipeOnDataNodes(request)).getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format("Failed to start pipe instance [%s] on data nodes",
pipeName));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
@@ -131,18 +121,10 @@ public class StartPipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
@Override
protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
- throws PipeManagementException {
+ throws PipeManagementException, IOException {
LOGGER.info("StartPipeProcedureV2: rollbackFromOperateOnDataNodes({})",
pipeName);
- final TOperatePipeOnDataNodeReq request =
- new TOperatePipeOnDataNodeReq()
- .setPipeName(pipeName)
- .setOperation((byte) PipeTaskOperation.STOP_PIPE.ordinal());
- if
(RpcUtils.squashResponseStatusList(env.operatePipeOnDataNodes(request)).getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format("Failed to rollback from start on data nodes for task
[%s]", pipeName));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StopPipeProcedureV2.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StopPipeProcedureV2.java
index 6726cfd5aeb..8f7f24b6148 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StopPipeProcedureV2.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/task/StopPipeProcedureV2.java
@@ -24,11 +24,8 @@ import
org.apache.iotdb.confignode.persistence.pipe.PipeTaskOperation;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
import org.apache.iotdb.consensus.common.response.ConsensusWriteResponse;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.iotdb.pipe.api.exception.PipeManagementException;
-import org.apache.iotdb.rpc.RpcUtils;
-import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import org.slf4j.Logger;
@@ -90,18 +87,11 @@ public class StopPipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
}
@Override
- void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws
PipeManagementException {
+ void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+ throws PipeManagementException, IOException {
LOGGER.info("StopPipeProcedureV2: executeFromOperateOnDataNodes({})",
pipeName);
- final TOperatePipeOnDataNodeReq request =
- new TOperatePipeOnDataNodeReq()
- .setPipeName(pipeName)
- .setOperation((byte) PipeTaskOperation.STOP_PIPE.ordinal());
- if
(RpcUtils.squashResponseStatusList(env.operatePipeOnDataNodes(request)).getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format("Failed to stop pipe instance [%s] on data nodes",
pipeName));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
@@ -131,18 +121,10 @@ public class StopPipeProcedureV2 extends
AbstractOperatePipeProcedureV2 {
@Override
protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
- throws PipeManagementException {
+ throws PipeManagementException, IOException {
LOGGER.info("StopPipeProcedureV2: rollbackFromOperateOnDataNodes({})",
pipeName);
- final TOperatePipeOnDataNodeReq request =
- new TOperatePipeOnDataNodeReq()
- .setPipeName(pipeName)
- .setOperation((byte) PipeTaskOperation.START_PIPE.ordinal());
- if
(RpcUtils.squashResponseStatusList(env.operatePipeOnDataNodes(request)).getCode()
- != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- throw new PipeManagementException(
- String.format("Failed to rollback from stop on data nodes for task
[%s]", pipeName));
- }
+ pushPipeMetaToDataNodes(env);
}
@Override
diff --git
a/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PipeInfoTest.java
b/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PipeInfoTest.java
index 0246759c98c..a547fdde521 100644
---
a/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PipeInfoTest.java
+++
b/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PipeInfoTest.java
@@ -73,6 +73,7 @@ public class PipeInfoTest {
processorAttributes.put("processor",
"org.apache.iotdb.pipe.processor.SDTFilterProcessor");
connectorAttributes.put("connector",
"org.apache.iotdb.pipe.protocal.ThriftTransporter");
PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(0, 1);
+ pipeTaskMeta.trackException(true, "someError");
Map<TConsensusGroupId, PipeTaskMeta> pipeTasks = new HashMap<>();
pipeTasks.put(new TConsensusGroupId(DataRegion, 1), pipeTaskMeta);
PipeStaticMeta pipeStaticMeta =
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMeta.java
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMeta.java
index f75237bbdc6..1fa610cc35c 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMeta.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMeta.java
@@ -65,8 +65,14 @@ public class PipeMeta {
}
public static PipeMeta deserialize(FileInputStream fileInputStream) throws
IOException {
- PipeStaticMeta staticMeta = PipeStaticMeta.deserialize(fileInputStream);
- PipeRuntimeMeta runtimeMeta = PipeRuntimeMeta.deserialize(fileInputStream);
+ final PipeStaticMeta staticMeta =
PipeStaticMeta.deserialize(fileInputStream);
+ final PipeRuntimeMeta runtimeMeta =
PipeRuntimeMeta.deserialize(fileInputStream);
+ return new PipeMeta(staticMeta, runtimeMeta);
+ }
+
+ public static PipeMeta deserialize(ByteBuffer byteBuffer) {
+ final PipeStaticMeta staticMeta = PipeStaticMeta.deserialize(byteBuffer);
+ final PipeRuntimeMeta runtimeMeta =
PipeRuntimeMeta.deserialize(byteBuffer);
return new PipeMeta(staticMeta, runtimeMeta);
}
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMetaKeeper.java
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMetaKeeper.java
index b45286abaf9..d5f0afb9188 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMetaKeeper.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeMetaKeeper.java
@@ -51,6 +51,10 @@ public class PipeMetaKeeper {
return pipeNameToPipeMetaMap.containsKey(pipeName);
}
+ public Iterable<PipeMeta> getPipeMetaList() {
+ return pipeNameToPipeMetaMap.values();
+ }
+
public void clear() {
this.pipeNameToPipeMetaMap.clear();
}
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeRuntimeMeta.java
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeRuntimeMeta.java
index 73551416e36..683177573ad 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeRuntimeMeta.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeRuntimeMeta.java
@@ -28,8 +28,6 @@ import java.io.DataOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.nio.ByteBuffer;
-import java.util.LinkedList;
-import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
@@ -38,18 +36,15 @@ import java.util.concurrent.atomic.AtomicReference;
public class PipeRuntimeMeta {
private final AtomicReference<PipeStatus> status;
- private final List<String> exceptionMessages;
private final Map<TConsensusGroupId, PipeTaskMeta>
consensusGroupIdToTaskMetaMap;
public PipeRuntimeMeta() {
status = new AtomicReference<>(PipeStatus.STOPPED);
- exceptionMessages = new LinkedList<>();
consensusGroupIdToTaskMetaMap = new ConcurrentHashMap<>();
}
public PipeRuntimeMeta(Map<TConsensusGroupId, PipeTaskMeta>
consensusGroupIdToTaskMetaMap) {
status = new AtomicReference<>(PipeStatus.STOPPED);
- exceptionMessages = new LinkedList<>();
this.consensusGroupIdToTaskMetaMap = consensusGroupIdToTaskMetaMap;
}
@@ -57,10 +52,6 @@ public class PipeRuntimeMeta {
return status;
}
- public List<String> getExceptionMessages() {
- return exceptionMessages;
- }
-
public Map<TConsensusGroupId, PipeTaskMeta>
getConsensusGroupIdToTaskMetaMap() {
return consensusGroupIdToTaskMetaMap;
}
@@ -75,8 +66,6 @@ public class PipeRuntimeMeta {
public void serialize(DataOutputStream outputStream) throws IOException {
ReadWriteIOUtils.write(status.get().getType(), outputStream);
- // ignore exception messages
-
ReadWriteIOUtils.write(consensusGroupIdToTaskMetaMap.size(), outputStream);
for (Map.Entry<TConsensusGroupId, PipeTaskMeta> entry :
consensusGroupIdToTaskMetaMap.entrySet()) {
@@ -90,13 +79,11 @@ public class PipeRuntimeMeta {
ByteBuffer.wrap(ReadWriteIOUtils.readBytesWithSelfDescriptionLength(inputStream)));
}
- public static PipeRuntimeMeta deserialize(ByteBuffer byteBuffer) throws
IOException {
+ public static PipeRuntimeMeta deserialize(ByteBuffer byteBuffer) {
final PipeRuntimeMeta pipeRuntimeMeta = new PipeRuntimeMeta();
pipeRuntimeMeta.status.set(PipeStatus.getPipeStatus(ReadWriteIOUtils.readByte(byteBuffer)));
- // ignore exception messages
-
final int size = ReadWriteIOUtils.readInt(byteBuffer);
for (int i = 0; i < size; ++i) {
pipeRuntimeMeta.consensusGroupIdToTaskMetaMap.put(
@@ -118,13 +105,12 @@ public class PipeRuntimeMeta {
}
PipeRuntimeMeta that = (PipeRuntimeMeta) o;
return Objects.equals(status.get().getType(), that.status.get().getType())
- && exceptionMessages.equals(that.exceptionMessages)
&&
consensusGroupIdToTaskMetaMap.equals(that.consensusGroupIdToTaskMetaMap);
}
@Override
public int hashCode() {
- return Objects.hash(status, exceptionMessages,
consensusGroupIdToTaskMetaMap);
+ return Objects.hash(status, consensusGroupIdToTaskMetaMap);
}
@Override
@@ -132,8 +118,6 @@ public class PipeRuntimeMeta {
return "PipeRuntimeMeta{"
+ "status="
+ status
- + ", exceptionMessages="
- + exceptionMessages
+ ", consensusGroupIdToTaskMetaMap="
+ consensusGroupIdToTaskMetaMap
+ '}';
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeStaticMeta.java
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeStaticMeta.java
index 3b9c94c94ce..14185dd8a5e 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeStaticMeta.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeStaticMeta.java
@@ -34,10 +34,6 @@ public class PipeStaticMeta {
private String pipeName;
private long createTime;
- private Map<String, String> collectorAttributes = new HashMap<>();
- private Map<String, String> processorAttributes = new HashMap<>();
- private Map<String, String> connectorAttributes = new HashMap<>();
-
private PipeParameters collectorParameters;
private PipeParameters processorParameters;
private PipeParameters connectorParameters;
@@ -52,9 +48,6 @@ public class PipeStaticMeta {
Map<String, String> connectorAttributes) {
this.pipeName = pipeName.toUpperCase();
this.createTime = createTime;
- this.collectorAttributes = collectorAttributes;
- this.processorAttributes = processorAttributes;
- this.connectorAttributes = connectorAttributes;
collectorParameters = new PipeParameters(collectorAttributes);
processorParameters = new PipeParameters(processorAttributes);
connectorParameters = new PipeParameters(connectorAttributes);
@@ -91,18 +84,18 @@ public class PipeStaticMeta {
ReadWriteIOUtils.write(pipeName, outputStream);
ReadWriteIOUtils.write(createTime, outputStream);
- outputStream.writeInt(collectorAttributes.size());
- for (Map.Entry<String, String> entry : collectorAttributes.entrySet()) {
+ outputStream.writeInt(collectorParameters.getAttribute().size());
+ for (Map.Entry<String, String> entry :
collectorParameters.getAttribute().entrySet()) {
ReadWriteIOUtils.write(entry.getKey(), outputStream);
ReadWriteIOUtils.write(entry.getValue(), outputStream);
}
- outputStream.writeInt(processorAttributes.size());
- for (Map.Entry<String, String> entry : processorAttributes.entrySet()) {
+ outputStream.writeInt(processorParameters.getAttribute().size());
+ for (Map.Entry<String, String> entry :
processorParameters.getAttribute().entrySet()) {
ReadWriteIOUtils.write(entry.getKey(), outputStream);
ReadWriteIOUtils.write(entry.getValue(), outputStream);
}
- outputStream.writeInt(connectorAttributes.size());
- for (Map.Entry<String, String> entry : connectorAttributes.entrySet()) {
+ outputStream.writeInt(connectorParameters.getAttribute().size());
+ for (Map.Entry<String, String> entry :
connectorParameters.getAttribute().entrySet()) {
ReadWriteIOUtils.write(entry.getKey(), outputStream);
ReadWriteIOUtils.write(entry.getValue(), outputStream);
}
@@ -119,26 +112,32 @@ public class PipeStaticMeta {
pipeStaticMeta.pipeName = ReadWriteIOUtils.readString(byteBuffer);
pipeStaticMeta.createTime = ReadWriteIOUtils.readLong(byteBuffer);
+ pipeStaticMeta.collectorParameters = new PipeParameters(new HashMap<>());
+ pipeStaticMeta.processorParameters = new PipeParameters(new HashMap<>());
+ pipeStaticMeta.connectorParameters = new PipeParameters(new HashMap<>());
+
int size = byteBuffer.getInt();
for (int i = 0; i < size; ++i) {
- pipeStaticMeta.collectorAttributes.put(
- ReadWriteIOUtils.readString(byteBuffer),
ReadWriteIOUtils.readString(byteBuffer));
+ pipeStaticMeta
+ .collectorParameters
+ .getAttribute()
+ .put(ReadWriteIOUtils.readString(byteBuffer),
ReadWriteIOUtils.readString(byteBuffer));
}
size = byteBuffer.getInt();
for (int i = 0; i < size; ++i) {
- pipeStaticMeta.processorAttributes.put(
- ReadWriteIOUtils.readString(byteBuffer),
ReadWriteIOUtils.readString(byteBuffer));
+ pipeStaticMeta
+ .processorParameters
+ .getAttribute()
+ .put(ReadWriteIOUtils.readString(byteBuffer),
ReadWriteIOUtils.readString(byteBuffer));
}
size = byteBuffer.getInt();
for (int i = 0; i < size; ++i) {
- pipeStaticMeta.connectorAttributes.put(
- ReadWriteIOUtils.readString(byteBuffer),
ReadWriteIOUtils.readString(byteBuffer));
+ pipeStaticMeta
+ .connectorParameters
+ .getAttribute()
+ .put(ReadWriteIOUtils.readString(byteBuffer),
ReadWriteIOUtils.readString(byteBuffer));
}
- pipeStaticMeta.collectorParameters = new
PipeParameters(pipeStaticMeta.collectorAttributes);
- pipeStaticMeta.processorParameters = new
PipeParameters(pipeStaticMeta.processorAttributes);
- pipeStaticMeta.connectorParameters = new
PipeParameters(pipeStaticMeta.connectorAttributes);
-
return pipeStaticMeta;
}
@@ -153,9 +152,9 @@ public class PipeStaticMeta {
PipeStaticMeta that = (PipeStaticMeta) obj;
return pipeName.equals(that.pipeName)
&& createTime == that.createTime
- && collectorAttributes.equals(that.collectorAttributes)
- && processorAttributes.equals(that.processorAttributes)
- && connectorAttributes.equals(that.connectorAttributes);
+ && collectorParameters.equals(that.collectorParameters)
+ && processorParameters.equals(that.processorParameters)
+ && connectorParameters.equals(that.connectorParameters);
}
@Override
@@ -171,12 +170,12 @@ public class PipeStaticMeta {
+ '\''
+ ", createTime="
+ createTime
- + ", collectorAttributes="
- + collectorAttributes
- + ", processorAttributes="
- + processorAttributes
- + ", connectorAttributes="
- + connectorAttributes
+ + ", collectorParameters="
+ + collectorParameters.getAttribute()
+ + ", processorParameters="
+ + processorParameters.getAttribute()
+ + ", connectorParameters="
+ + connectorParameters.getAttribute()
+ '}';
}
}
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeTaskMeta.java
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeTaskMeta.java
index 258a7a73503..4cd9402ba3d 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeTaskMeta.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/pipe/task/meta/PipeTaskMeta.java
@@ -19,12 +19,19 @@
package org.apache.iotdb.commons.pipe.task.meta;
+import org.apache.iotdb.pipe.api.exception.PipeRuntimeCriticalException;
+import org.apache.iotdb.pipe.api.exception.PipeRuntimeException;
+import org.apache.iotdb.pipe.api.exception.PipeRuntimeNonCriticalException;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import java.io.DataOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.Objects;
+import java.util.Queue;
+import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
@@ -33,6 +40,7 @@ public class PipeTaskMeta {
// TODO: replace it with consensus index
private final AtomicLong index = new AtomicLong(0L);
private final AtomicInteger regionLeader = new AtomicInteger(0);
+ private final Queue<PipeRuntimeException> exceptionMessages = new
ConcurrentLinkedQueue<>();
private PipeTaskMeta() {}
@@ -49,6 +57,17 @@ public class PipeTaskMeta {
return regionLeader.get();
}
+ public Iterable<PipeRuntimeException> getExceptionMessages() {
+ return exceptionMessages;
+ }
+
+ public void trackException(boolean critical, String message) {
+ exceptionMessages.add(
+ critical
+ ? new PipeRuntimeCriticalException(message)
+ : new PipeRuntimeNonCriticalException(message));
+ }
+
public void setIndex(long index) {
this.index.set(index);
}
@@ -60,12 +79,27 @@ public class PipeTaskMeta {
public void serialize(DataOutputStream outputStream) throws IOException {
ReadWriteIOUtils.write(index.get(), outputStream);
ReadWriteIOUtils.write(regionLeader.get(), outputStream);
+ ReadWriteIOUtils.write(exceptionMessages.size(), outputStream);
+ for (final PipeRuntimeException exceptionMessage : exceptionMessages) {
+ ReadWriteIOUtils.write(
+ exceptionMessage instanceof PipeRuntimeCriticalException,
outputStream);
+ ReadWriteIOUtils.write(exceptionMessage.getMessage(), outputStream);
+ }
}
public static PipeTaskMeta deserialize(ByteBuffer byteBuffer) {
final PipeTaskMeta PipeTaskMeta = new PipeTaskMeta();
PipeTaskMeta.index.set(ReadWriteIOUtils.readLong(byteBuffer));
PipeTaskMeta.regionLeader.set(ReadWriteIOUtils.readInt(byteBuffer));
+ final int size = ReadWriteIOUtils.readInt(byteBuffer);
+ for (int i = 0; i < size; ++i) {
+ final boolean critical = ReadWriteIOUtils.readBool(byteBuffer);
+ final String message = ReadWriteIOUtils.readString(byteBuffer);
+ PipeTaskMeta.exceptionMessages.add(
+ critical
+ ? new PipeRuntimeCriticalException(message)
+ : new PipeRuntimeNonCriticalException(message));
+ }
return PipeTaskMeta;
}
@@ -83,16 +117,27 @@ public class PipeTaskMeta {
return false;
}
PipeTaskMeta that = (PipeTaskMeta) obj;
- return index.get() == that.index.get() && regionLeader.get() ==
that.regionLeader.get();
+ return index.get() == that.index.get()
+ && regionLeader.get() == that.regionLeader.get()
+ && Arrays.equals(exceptionMessages.toArray(),
that.exceptionMessages.toArray());
}
@Override
public int hashCode() {
- return (int) (index.get() * 31 + regionLeader.get());
+ return Objects.hash(index, regionLeader, exceptionMessages);
}
@Override
public String toString() {
- return "PipeTask{" + "index='" + index + '\'' + ", regionLeader='" +
regionLeader + '\'' + '}';
+ return "PipeTask{"
+ + "index='"
+ + index
+ + '\''
+ + ", regionLeader='"
+ + regionLeader
+ + '\''
+ + ", exceptionMessages="
+ + exceptionMessages
+ + '}';
}
}
diff --git
a/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/PipeParameters.java
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/PipeParameters.java
index 18dbfdd4ab7..1dd852d9b36 100644
---
a/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/PipeParameters.java
+++
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/PipeParameters.java
@@ -109,4 +109,26 @@ public class PipeParameters {
String value = attributes.get(key);
return value == null ? defaultValue : Double.parseDouble(value);
}
+
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj) {
+ return true;
+ }
+ if (obj == null || getClass() != obj.getClass()) {
+ return false;
+ }
+ PipeParameters that = (PipeParameters) obj;
+ return attributes.equals(that.attributes);
+ }
+
+ @Override
+ public int hashCode() {
+ return attributes.hashCode();
+ }
+
+ @Override
+ public String toString() {
+ return attributes.toString();
+ }
}
diff --git
a/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeCriticalException.java
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeCriticalException.java
new file mode 100644
index 00000000000..fbe73118385
--- /dev/null
+++
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeCriticalException.java
@@ -0,0 +1,40 @@
+/*
+ * 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.pipe.api.exception;
+
+import java.util.Objects;
+
+public class PipeRuntimeCriticalException extends PipeRuntimeException {
+
+ public PipeRuntimeCriticalException(String message) {
+ super(message);
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ return obj instanceof PipeRuntimeCriticalException
+ && Objects.equals(getMessage(), ((PipeRuntimeCriticalException)
obj).getMessage());
+ }
+
+ @Override
+ public int hashCode() {
+ return getMessage().hashCode();
+ }
+}
diff --git
a/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeException.java
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeException.java
new file mode 100644
index 00000000000..9b20dff7bff
--- /dev/null
+++
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeException.java
@@ -0,0 +1,40 @@
+/*
+ * 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.pipe.api.exception;
+
+import java.util.Objects;
+
+public abstract class PipeRuntimeException extends PipeException {
+
+ public PipeRuntimeException(String message) {
+ super(message);
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ return obj instanceof PipeRuntimeException
+ && Objects.equals(getMessage(), ((PipeRuntimeException)
obj).getMessage());
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(getMessage());
+ }
+}
diff --git
a/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeNonCriticalException.java
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeNonCriticalException.java
new file mode 100644
index 00000000000..6fa63b93cc1
--- /dev/null
+++
b/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeRuntimeNonCriticalException.java
@@ -0,0 +1,40 @@
+/*
+ * 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.pipe.api.exception;
+
+import java.util.Objects;
+
+public class PipeRuntimeNonCriticalException extends PipeRuntimeException {
+
+ public PipeRuntimeNonCriticalException(String message) {
+ super(message);
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ return obj instanceof PipeRuntimeNonCriticalException
+ && Objects.equals(getMessage(), ((PipeRuntimeNonCriticalException)
obj).getMessage());
+ }
+
+ @Override
+ public int hashCode() {
+ return getMessage().hashCode();
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
index 13690778a09..a0675c0fb63 100644
---
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -44,7 +44,6 @@ import
org.apache.iotdb.commons.pipe.plugin.meta.PipePluginMeta;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.commons.service.metric.enums.Metric;
import org.apache.iotdb.commons.service.metric.enums.Tag;
-import org.apache.iotdb.commons.sync.pipe.SyncOperation;
import org.apache.iotdb.commons.trigger.TriggerInformation;
import org.apache.iotdb.commons.udf.UDFInformation;
import org.apache.iotdb.commons.udf.service.UDFManagementService;
@@ -137,7 +136,6 @@ import
org.apache.iotdb.mpp.rpc.thrift.TCountPathsUsingTemplateResp;
import org.apache.iotdb.mpp.rpc.thrift.TCreateDataRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateFunctionInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreatePeerReq;
-import org.apache.iotdb.mpp.rpc.thrift.TCreatePipeOnDataNodeReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreatePipePluginInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateSchemaRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TCreateTriggerInstanceReq;
@@ -166,7 +164,7 @@ import org.apache.iotdb.mpp.rpc.thrift.TLoadCommandReq;
import org.apache.iotdb.mpp.rpc.thrift.TLoadResp;
import org.apache.iotdb.mpp.rpc.thrift.TLoadSample;
import org.apache.iotdb.mpp.rpc.thrift.TMaintainPeerReq;
-import org.apache.iotdb.mpp.rpc.thrift.TOperatePipeOnDataNodeReq;
+import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaReq;
import org.apache.iotdb.mpp.rpc.thrift.TRegionLeaderChangeReq;
import org.apache.iotdb.mpp.rpc.thrift.TRegionRouteReq;
import org.apache.iotdb.mpp.rpc.thrift.TRollbackSchemaBlackListReq;
@@ -802,6 +800,11 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
return resp;
}
+ @Override
+ public TSStatus pushPipeMeta(TPushPipeMetaReq req) throws TException {
+ return null;
+ }
+
private TSStatus executeInternalSchemaTask(
List<TConsensusGroupId> consensusGroupIdList,
Function<TConsensusGroupId, TSStatus> executeOnOneRegion) {
@@ -822,28 +825,6 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
}
}
- @Override
- public TSStatus createPipeOnDataNode(TCreatePipeOnDataNodeReq req) {
- throw new NotImplementedException("TODO: createPipeOnDataNode");
- }
-
- @Override
- public TSStatus operatePipeOnDataNode(TOperatePipeOnDataNodeReq req) {
- try {
- switch (SyncOperation.values()[req.getOperation()]) {
- case START_PIPE:
- case STOP_PIPE:
- case DROP_PIPE:
- throw new NotImplementedException("TODO: operatePipeOnDataNode");
- default:
- return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
- .setMessage("Unsupported operation.");
- }
- } catch (Exception e) {
- return new
TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode()).setMessage(e.getMessage());
- }
- }
-
@Override
public TSStatus executeCQ(TExecuteCQ req) {
diff --git a/thrift/src/main/thrift/datanode.thrift
b/thrift/src/main/thrift/datanode.thrift
index a9017c2c3ce..53817e266a2 100644
--- a/thrift/src/main/thrift/datanode.thrift
+++ b/thrift/src/main/thrift/datanode.thrift
@@ -373,14 +373,8 @@ struct TCheckTimeSeriesExistenceResp{
2: optional bool exists
}
-struct TCreatePipeOnDataNodeReq{
- 1: required binary pipeMeta
-}
-
-struct TOperatePipeOnDataNodeReq {
- 1: required string pipeName
- // ordinal of {@linkplain SyncOperation}
- 2: required i8 operation
+struct TPushPipeMetaReq {
+ 1: required list<binary> pipeMetas
}
// ====================================================
@@ -761,14 +755,9 @@ service IDataNodeRPCService {
TCheckTimeSeriesExistenceResp
checkTimeSeriesExistence(TCheckTimeSeriesExistenceReq req)
/**
- * Create PIPE on DataNode
- */
- common.TSStatus createPipeOnDataNode(TCreatePipeOnDataNodeReq req)
-
- /**
- * Start, stop or drop PIPE on DataNode
+ * Send pipeMetas to DataNodes, for synchronization
*/
- common.TSStatus operatePipeOnDataNode(TOperatePipeOnDataNodeReq req)
+ common.TSStatus pushPipeMeta(TPushPipeMetaReq req)
/**
* Execute CQ on DataNode