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

Reply via email to