This is an automated email from the ASF dual-hosted git repository.

jt2594838 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 9ebbaaa458d Fix snapshot pipe auto-drop when some DataNodes fail to 
initialize (#18512)
9ebbaaa458d is described below

commit 9ebbaaa458d8c9e23396b40708d922f5ed570b14
Author: Zhenyu Luo <[email protected]>
AuthorDate: Tue Aug 25 16:01:52 2026 +0800

    Fix snapshot pipe auto-drop when some DataNodes fail to initialize (#18512)
    
    * Fix snapshot pipe auto-drop when some DataNodes fail to initialize
    
    * Fix pipe snapshot auto-drop by completing only after all required 
DataRegions report finished
    
    * Fix IllegalPathException catch in PipeDataNodeTaskAgent after removing DN 
completion boolean
    
    * Fix unhandled IOException in PipeHeartbeatParserTest helper
    
    * Fix pipe TsFile decomposition to filter table-model mods deletions
    
    When a TsFile is split into tablets for pipe transfer, table-model 
deletions recorded in mods2 were not matched because ModsOperationUtil queried 
the modification pattern tree with an IDeviceID directly. Convert the device 
and measurement to a full path before lookup, matching ModificationUtils 
behavior, so deleted rows are filtered on the sender side.
---
 .../heartbeat/DataNodeHeartbeatHandler.java        |   3 +-
 .../iotdb/confignode/manager/ConfigManager.java    |   3 +-
 .../runtime/PipeRuntimeCoordinator.java            |   7 +-
 .../runtime/heartbeat/PipeHeartbeat.java           |  46 +++++--
 .../runtime/heartbeat/PipeHeartbeatParser.java     |  72 ++++++-----
 .../runtime/heartbeat/PipeHeartbeatScheduler.java  |   3 +-
 .../confignode/persistence/pipe/PipeTaskInfo.java  |   4 -
 .../runtime/heartbeat/PipeHeartbeatParserTest.java | 135 ++++++++++++++++++++-
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  |  99 +++++++--------
 .../tsfile/parser/util/ModsOperationUtil.java      |  17 ++-
 .../task/meta/PipeTemporaryMetaInCoordinator.java  |  22 ++--
 .../thrift-commons/src/main/thrift/common.thrift   |   7 ++
 .../src/main/thrift/datanode.thrift                |   1 +
 13 files changed, 305 insertions(+), 114 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/heartbeat/DataNodeHeartbeatHandler.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/heartbeat/DataNodeHeartbeatHandler.java
index 9c7810dabe2..14e1dcbb86f 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/heartbeat/DataNodeHeartbeatHandler.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/heartbeat/DataNodeHeartbeatHandler.java
@@ -192,7 +192,8 @@ public class DataNodeHeartbeatHandler implements 
AsyncMethodCallback<TDataNodeHe
           heartbeatResp.getPipeRemainingEventCountList(),
           heartbeatResp.getPipeRemainingTimeList(),
           heartbeatResp.getPipeDegradedStatusList(),
-          heartbeatResp.getPipeRecentFailureList());
+          heartbeatResp.getPipeRecentFailureList(),
+          heartbeatResp.getPipeCompletedDataRegionList());
     }
   }
 
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
index c0276fd2cd5..7ff21974540 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
@@ -3339,7 +3339,8 @@ public class ConfigManager implements IManager {
             resp.getPipeRemainingEventCountList(),
             resp.getPipeRemainingTimeList(),
             resp.getPipeDegradedStatusList(),
-            resp.getPipeRecentFailureList());
+            resp.getPipeRecentFailureList(),
+            resp.getPipeCompletedDataRegionList());
     return StatusUtils.OK;
   }
 
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/PipeRuntimeCoordinator.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/PipeRuntimeCoordinator.java
index ec00adcd302..cb0da3e0948 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/PipeRuntimeCoordinator.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/PipeRuntimeCoordinator.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.confignode.manager.pipe.coordinator.runtime;
 
+import org.apache.iotdb.common.rpc.thrift.TPipeCompletedDataRegion;
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.commons.concurrent.ThreadName;
 import org.apache.iotdb.confignode.manager.ConfigManager;
@@ -98,7 +99,8 @@ public class PipeRuntimeCoordinator implements 
IClusterStatusSubscriber {
       /* @Nullable */ final List<Long> pipeRemainingEventCountListFromAgent,
       /* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent,
       /* @Nullable */ final List<Integer> pipeDegradedStatusListFromAgent,
-      /* @Nullable */ final List<Map<String, Long>> 
pipeRecentFailureListFromAgent) {
+      /* @Nullable */ final List<Map<String, Long>> 
pipeRecentFailureListFromAgent,
+      /* @Nullable */ final List<TPipeCompletedDataRegion> 
pipeCompletedDataRegionListFromAgent) {
     pipeHeartbeatScheduler.parseHeartbeat(
         dataNodeId,
         new PipeHeartbeat(
@@ -107,6 +109,7 @@ public class PipeRuntimeCoordinator implements 
IClusterStatusSubscriber {
             pipeRemainingEventCountListFromAgent,
             pipeRemainingTimeListFromAgent,
             pipeDegradedStatusListFromAgent,
-            pipeRecentFailureListFromAgent));
+            pipeRecentFailureListFromAgent,
+            pipeCompletedDataRegionListFromAgent));
   }
 }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeat.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeat.java
index 7aa75b2d78d..dd31751bb17 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeat.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeat.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.confignode.manager.pipe.coordinator.runtime.heartbeat;
 
+import org.apache.iotdb.common.rpc.thrift.TPipeCompletedDataRegion;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMeta;
@@ -32,11 +33,11 @@ import java.util.Objects;
 
 public class PipeHeartbeat {
   private final Map<PipeStaticMeta, PipeMeta> pipeMetaMap = new HashMap<>();
-  private final Map<PipeStaticMeta, Boolean> isCompletedMap = new HashMap<>();
   private final Map<PipeStaticMeta, Long> remainingEventCountMap = new 
HashMap<>();
   private final Map<PipeStaticMeta, Double> remainingTimeMap = new HashMap<>();
   private final Map<PipeStaticMeta, Boolean> isDegradedMap = new HashMap<>();
   private final Map<PipeStaticMeta, Map<String, Long>> recentFailuresMap = new 
HashMap<>();
+  private final Map<PipeStaticMeta, List<Integer>> completedDataRegionIdsMap = 
new HashMap<>();
 
   public PipeHeartbeat(
       final List<ByteBuffer> pipeMetaByteBufferListFromAgent,
@@ -50,6 +51,7 @@ public class PipeHeartbeat {
         pipeRemainingEventCountListFromAgent,
         pipeRemainingTimeListFromAgent,
         pipeDegradedStatusListFromAgent,
+        null,
         null);
   }
 
@@ -60,6 +62,24 @@ public class PipeHeartbeat {
       /* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent,
       /* @Nullable */ final List<Integer> pipeDegradedStatusListFromAgent,
       /* @Nullable */ final List<Map<String, Long>> 
pipeRecentFailureListFromAgent) {
+    this(
+        pipeMetaByteBufferListFromAgent,
+        pipeCompletedListFromAgent,
+        pipeRemainingEventCountListFromAgent,
+        pipeRemainingTimeListFromAgent,
+        pipeDegradedStatusListFromAgent,
+        pipeRecentFailureListFromAgent,
+        null);
+  }
+
+  public PipeHeartbeat(
+      final List<ByteBuffer> pipeMetaByteBufferListFromAgent,
+      /* @Nullable */ final List<Boolean> pipeCompletedListFromAgent,
+      /* @Nullable */ final List<Long> pipeRemainingEventCountListFromAgent,
+      /* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent,
+      /* @Nullable */ final List<Integer> pipeDegradedStatusListFromAgent,
+      /* @Nullable */ final List<Map<String, Long>> 
pipeRecentFailureListFromAgent,
+      /* @Nullable */ final List<TPipeCompletedDataRegion> 
completedDataRegionListFromAgent) {
     // Shall not reach here, just in case
     if (Objects.isNull(pipeMetaByteBufferListFromAgent)) {
       return;
@@ -68,11 +88,6 @@ public class PipeHeartbeat {
       final PipeMeta pipeMeta =
           
PipeMeta.deserialize4TaskAgent(pipeMetaByteBufferListFromAgent.get(i));
       pipeMetaMap.put(pipeMeta.getStaticMeta(), pipeMeta);
-      isCompletedMap.put(
-          pipeMeta.getStaticMeta(),
-          Objects.nonNull(pipeCompletedListFromAgent)
-              && i < pipeCompletedListFromAgent.size()
-              && pipeCompletedListFromAgent.get(i));
       // If remaining event count & remaining time can not be got, it implies 
that the heartbeat is
       // from an ancient version of DataNode. Here we guarantee that "0" will 
not affect both of
       // the final results and namely these dataNodes are omitted in 
calculation.
@@ -102,6 +117,17 @@ public class PipeHeartbeat {
                   && Objects.nonNull(pipeRecentFailureListFromAgent.get(i))
               ? new HashMap<>(pipeRecentFailureListFromAgent.get(i))
               : Collections.emptyMap());
+      if (completedDataRegionListFromAgent != null) {
+        for (final TPipeCompletedDataRegion completedDataRegion :
+            completedDataRegionListFromAgent) {
+          if 
(pipeMeta.getStaticMeta().getPipeName().equals(completedDataRegion.getPipeName())
+              && pipeMeta.getStaticMeta().getCreationTime()
+                  == completedDataRegion.getCreationTime()) {
+            completedDataRegionIdsMap.put(
+                pipeMeta.getStaticMeta(), 
completedDataRegion.getCompletedDataRegionIds());
+          }
+        }
+      }
     }
   }
 
@@ -113,8 +139,12 @@ public class PipeHeartbeat {
     return pipeMetaMap.get(pipeStaticMeta);
   }
 
-  public Boolean isCompleted(final PipeStaticMeta pipeStaticMeta) {
-    return isCompletedMap.get(pipeStaticMeta);
+  public List<Integer> getCompletedDataRegionIds(final PipeStaticMeta 
pipeStaticMeta) {
+    return completedDataRegionIdsMap.getOrDefault(pipeStaticMeta, 
Collections.emptyList());
+  }
+
+  public boolean hasCompletedDataRegionReport(final PipeStaticMeta 
pipeStaticMeta) {
+    return completedDataRegionIdsMap.containsKey(pipeStaticMeta);
   }
 
   public Long getRemainingEventCount(final PipeStaticMeta pipeStaticMeta) {
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
index a8734469d0c..1149c8e4444 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
@@ -19,6 +19,8 @@
 
 package org.apache.iotdb.confignode.manager.pipe.coordinator.runtime.heartbeat;
 
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeCriticalException;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeException;
@@ -39,6 +41,7 @@ import 
org.apache.iotdb.confignode.persistence.pipe.PipeTaskInfo;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.HashSet;
 import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.atomic.AtomicBoolean;
@@ -156,41 +159,48 @@ public class PipeHeartbeatParser {
       final PipeTemporaryMetaInCoordinator temporaryMeta =
           (PipeTemporaryMetaInCoordinator) 
pipeMetaFromCoordinator.getTemporaryMeta();
 
-      // Remove completed pipes
-      final Boolean isPipeCompletedFromAgent = 
pipeHeartbeat.isCompleted(staticMeta);
-      if (Boolean.TRUE.equals(isPipeCompletedFromAgent)) {
+      // Aggregate completed DataRegion ids reported by DataNodes. Only the 
DataNodes that own the
+      // target region can report it, so the coordinator can compare the union 
against all required
+      // DataRegion ids without trusting any DataNode's single per-pipe 
completion boolean.
+      if (pipeHeartbeat.hasCompletedDataRegionReport(staticMeta)) {
+        for (final Integer completedDataRegionId :
+            pipeHeartbeat.getCompletedDataRegionIds(staticMeta)) {
+          temporaryMeta.markDataRegionCompleted(completedDataRegionId);
+        }
+      }
+
+      final Set<Integer> requiredDataRegionIds = new HashSet<>();
+      for (final Map.Entry<Integer, PipeTaskMeta> entry :
+          
pipeMetaFromCoordinator.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().entrySet())
 {
+        if (configManager
+            .getPartitionManager()
+            .isRegionGroupExists(
+                new TConsensusGroupId(TConsensusGroupType.DataRegion, 
entry.getKey()))) {
+          requiredDataRegionIds.add(entry.getKey());
+        }
+      }
 
-        temporaryMeta.markDataNodeCompleted(nodeId);
+      // Remove completed pipes only when every required DataRegion has been 
reported complete.
+      // Relying on the region-level reports (instead of the DataNode-level 
boolean) prevents a
+      // leader-change / task-creation failure from being treated as a 
successful snapshot transfer.
+      if (!requiredDataRegionIds.isEmpty()
+          && 
temporaryMeta.getCompletedDataRegionIds().containsAll(requiredDataRegionIds)) {
         PipeLogger.log(
             LOGGER::info,
-            
ManagerMessages.DETECTED_HISTORICAL_PIPE_COMPLETION_REPORT_FROM_DATANODE,
-            nodeId,
+            ManagerMessages.ALL_DATANODES_REPORTED_HISTORICAL_PIPE_COMPLETED,
             staticMeta.getPipeName(),
-            pipeHeartbeat.getRemainingEventCount(staticMeta),
-            pipeHeartbeat.getRemainingTime(staticMeta),
-            temporaryMeta.getCompletedDataNodeIds());
-
-        final Set<Integer> uncompletedDataNodeIds =
-            
configManager.getNodeManager().getRegisteredDataNodeLocations().keySet();
-        
uncompletedDataNodeIds.removeAll(temporaryMeta.getCompletedDataNodeIds());
-        if (uncompletedDataNodeIds.isEmpty()) {
-          PipeLogger.log(
-              LOGGER::info,
-              ManagerMessages.ALL_DATANODES_REPORTED_HISTORICAL_PIPE_COMPLETED,
-              staticMeta.getPipeName(),
-              temporaryMeta.getGlobalRemainingEvents(),
-              temporaryMeta.getGlobalRemainingTime(),
-              staticMeta);
-          pipeTaskInfo.get().removePipeMeta(staticMeta);
-          PipeLogger.log(
-              LOGGER::info,
-              
ManagerMessages.DETECTED_COMPLETION_OF_PIPE_STATIC_META_REMOVE_IT,
-              staticMeta.getPipeName(),
-              staticMeta);
-          needWriteConsensusOnConfigNodes.set(true);
-          needPushPipeMetaToDataNodes.set(true);
-          continue;
-        }
+            temporaryMeta.getGlobalRemainingEvents(),
+            temporaryMeta.getGlobalRemainingTime(),
+            staticMeta);
+        pipeTaskInfo.get().removePipeMeta(staticMeta);
+        PipeLogger.log(
+            LOGGER::info,
+            ManagerMessages.DETECTED_COMPLETION_OF_PIPE_STATIC_META_REMOVE_IT,
+            staticMeta.getPipeName(),
+            staticMeta);
+        needWriteConsensusOnConfigNodes.set(true);
+        needPushPipeMetaToDataNodes.set(true);
+        continue;
       }
 
       // Record statistics
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
index 209b08cff15..79d944ad87b 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatScheduler.java
@@ -117,7 +117,8 @@ public class PipeHeartbeatScheduler {
                         resp.getPipeRemainingEventCountList(),
                         resp.getPipeRemainingTimeList(),
                         resp.getPipeDegradedStatusList(),
-                        resp.getPipeRecentFailureList())));
+                        resp.getPipeRecentFailureList(),
+                        resp.getPipeCompletedDataRegionList())));
 
     // config node heartbeat
     try {
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
index 97856711050..2427e7f3329 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/pipe/PipeTaskInfo.java
@@ -34,7 +34,6 @@ import 
org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMeta;
-import 
org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInCoordinator;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeType;
 import org.apache.iotdb.commons.pipe.config.PipeConfig;
 import org.apache.iotdb.commons.pipe.config.constant.PipeProcessorConstant;
@@ -917,9 +916,6 @@ public class PipeTaskInfo implements SnapshotProcessor {
                               consensusGroupIdToTaskMetaMap
                                   .get(consensusGroupId.getId())
                                   .setLeaderNodeId(newLeader);
-                              // New region leader may contain un-transferred 
events
-                              ((PipeTemporaryMetaInCoordinator) 
pipeMeta.getTemporaryMeta())
-                                  .markDataNodeUncompleted(newLeader);
                             } else {
                               
consensusGroupIdToTaskMetaMap.remove(consensusGroupId.getId());
                             }
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
index c7d4e3b5d8d..95c29a56336 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
@@ -19,6 +19,8 @@
 
 package org.apache.iotdb.confignode.manager.pipe.coordinator.runtime.heartbeat;
 
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
+import org.apache.iotdb.common.rpc.thrift.TPipeCompletedDataRegion;
 import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeCriticalException;
@@ -33,6 +35,7 @@ import 
org.apache.iotdb.confignode.consensus.request.write.pipe.task.CreatePipeP
 import org.apache.iotdb.confignode.manager.ConfigManager;
 import org.apache.iotdb.confignode.manager.ProcedureManager;
 import org.apache.iotdb.confignode.manager.node.NodeManager;
+import org.apache.iotdb.confignode.manager.partition.PartitionManager;
 import org.apache.iotdb.confignode.manager.pipe.coordinator.PipeManager;
 import 
org.apache.iotdb.confignode.manager.pipe.coordinator.runtime.PipeRuntimeCoordinator;
 import 
org.apache.iotdb.confignode.manager.pipe.coordinator.task.PipeTaskCoordinator;
@@ -47,6 +50,7 @@ import org.mockito.Mockito;
 import java.lang.reflect.Field;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentHashMap;
@@ -360,6 +364,105 @@ public class PipeHeartbeatParserTest {
     
Assert.assertTrue(heartbeat.getRecentFailures(pipeMeta.getStaticMeta()).isEmpty());
   }
 
+  @Test
+  public void testParseHeartbeatDoesNotCompleteWhenRequiredDataRegionMissing() 
throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta();
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(1, pipeTaskInfo);
+    context.parser.parseHeartbeat(
+        1,
+        new PipeHeartbeat(
+            Collections.singletonList(pipeMeta.serialize()),
+            Collections.singletonList(true),
+            Collections.singletonList(0L),
+            Collections.singletonList(0d),
+            null,
+            null,
+            Collections.singletonList(
+                new TPipeCompletedDataRegion(
+                    pipeMeta.getStaticMeta().getPipeName(),
+                    pipeMeta.getStaticMeta().getCreationTime(),
+                    Collections.emptyList()))));
+
+    
Assert.assertTrue(getTemporaryMeta(pipeTaskInfo).getCompletedDataRegionIds().isEmpty());
+    Assert.assertNotNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+    verify(context.procedureManager, 
never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
+  }
+
+  @Test
+  public void testParseHeartbeatDoesNotTrustDataNodeBooleanForCompletion() 
throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta(1);
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(1, pipeTaskInfo);
+    // The DataNode's boolean is false, but the required DataRegion is 
reported complete. The
+    // coordinator should still complete the pipe because it no longer trusts 
the boolean.
+    context.parser.parseHeartbeat(
+        1, createPipeHeartbeatWithCompletedRegions(pipeMeta, false, 
Collections.singletonList(1)));
+
+    Assert.assertNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+  }
+
+  @Test
+  public void 
testParseHeartbeatCompletesOnlyAfterAllRequiredDataRegionsReported()
+      throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta(1, 2);
+    
pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().get(2).setLeaderNodeId(2);
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(2, pipeTaskInfo);
+
+    context.parser.parseHeartbeat(
+        1, createPipeHeartbeatWithCompletedRegions(pipeMeta, true, 
Collections.singletonList(1)));
+    Assert.assertNotNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+
+    context.parser.parseHeartbeat(
+        2, createPipeHeartbeatWithCompletedRegions(pipeMeta, true, 
Collections.singletonList(2)));
+    Assert.assertNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+    // After CN decides the pipe is complete, the next heartbeat round pushes 
the updated meta so
+    // DataNodes will drop their local pipe tasks.
+    verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, 
true);
+  }
+
+  @Test
+  public void testParseHeartbeatKeepsCompletedDataRegionAfterLeaderChange() 
throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta(1, 2);
+    
pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().get(2).setLeaderNodeId(2);
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(2, pipeTaskInfo);
+
+    // The old leader of region 1 reports it complete before the leader 
changes.
+    context.parser.parseHeartbeat(
+        1, createPipeHeartbeatWithCompletedRegions(pipeMeta, true, 
Collections.singletonList(1)));
+    Assert.assertNotNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+
+    // Region 1's leader moves to node 2, which only reports region 2. Region 
1's completion is
+    // still valid because its historical data was already transferred by the 
old leader.
+    
pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().get(1).setLeaderNodeId(2);
+    context.parser.parseHeartbeat(
+        2, createPipeHeartbeatWithCompletedRegions(pipeMeta, true, 
Collections.singletonList(2)));
+
+    Assert.assertNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+  }
+
   private ParserTestContext createParserTestContext(final int 
registeredDataNodeCount) {
     return createParserTestContext(registeredDataNodeCount, new 
PipeTaskInfo());
   }
@@ -373,6 +476,7 @@ public class PipeHeartbeatParserTest {
     final PipeRuntimeCoordinator pipeRuntimeCoordinator =
         Mockito.mock(PipeRuntimeCoordinator.class);
     final PipeTaskCoordinator pipeTaskCoordinator = 
Mockito.mock(PipeTaskCoordinator.class);
+    final PartitionManager partitionManager = 
Mockito.mock(PartitionManager.class);
     final ExecutorService procedureSubmitter = 
Mockito.mock(ExecutorService.class);
 
     when(configManager.getNodeManager()).thenReturn(nodeManager);
@@ -382,6 +486,8 @@ public class PipeHeartbeatParserTest {
     
when(pipeManager.getPipeRuntimeCoordinator()).thenReturn(pipeRuntimeCoordinator);
     when(pipeManager.getPipeTaskCoordinator()).thenReturn(pipeTaskCoordinator);
     
when(pipeRuntimeCoordinator.getProcedureSubmitter()).thenReturn(procedureSubmitter);
+    when(configManager.getPartitionManager()).thenReturn(partitionManager);
+    
when(partitionManager.isRegionGroupExists(any(TConsensusGroupId.class))).thenReturn(true);
     when(pipeTaskCoordinator.tryLock()).thenReturn(new 
AtomicReference<>(pipeTaskInfo));
     when(procedureManager.pipeHandleMetaChange(anyBoolean(), 
anyBoolean())).thenReturn(true);
     Mockito.doAnswer(
@@ -467,15 +573,38 @@ public class PipeHeartbeatParserTest {
   }
 
   private PipeMeta createPipeMeta() {
+    return createPipeMeta(1);
+  }
+
+  private PipeMeta createPipeMeta(final int... regionIds) {
     final PipeRuntimeMeta pipeRuntimeMeta = new PipeRuntimeMeta();
-    pipeRuntimeMeta
-        .getConsensusGroupId2TaskMetaMap()
-        .put(1, new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1));
+    for (final int regionId : regionIds) {
+      pipeRuntimeMeta
+          .getConsensusGroupId2TaskMetaMap()
+          .put(regionId, new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1));
+    }
     return new PipeMeta(
         new PipeStaticMeta("test_pipe", 1L, new HashMap<>(), new HashMap<>(), 
new HashMap<>()),
         pipeRuntimeMeta);
   }
 
+  private PipeHeartbeat createPipeHeartbeatWithCompletedRegions(
+      final PipeMeta pipeMeta, final boolean isCompleted, final List<Integer> 
completedRegionIds)
+      throws Exception {
+    return new PipeHeartbeat(
+        Collections.singletonList(pipeMeta.serialize()),
+        Collections.singletonList(isCompleted),
+        Collections.singletonList(0L),
+        Collections.singletonList(0d),
+        null,
+        null,
+        Collections.singletonList(
+            new TPipeCompletedDataRegion(
+                pipeMeta.getStaticMeta().getPipeName(),
+                pipeMeta.getStaticMeta().getCreationTime(),
+                completedRegionIds)));
+  }
+
   private PipeHeartbeat emptyHeartbeat() {
     return new PipeHeartbeat(Collections.emptyList(), null, null, null, null);
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
index 5241ea0c916..1f6416ad20f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.db.pipe.agent.task;
 
+import org.apache.iotdb.common.rpc.thrift.TPipeCompletedDataRegion;
 import org.apache.iotdb.common.rpc.thrift.TPipeHeartbeatResp;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.concurrent.IoTThreadFactory;
@@ -501,7 +502,7 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
                 
PipeConfig.getInstance().getPipeMetaReportMaxLogIntervalRounds(),
                 pipeMetaKeeper.getPipeMetaCount());
 
-    collectPipeMetaReport(logger, true).setTo(resp);
+    collectPipeMetaReport(logger).setTo(resp);
     PipeInsertionDataNodeListener.getInstance().listenToHeartbeat(true);
   }
 
@@ -523,17 +524,11 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
     LOGGER.debug(
         DataNodePipeMessages.RECEIVED_PIPE_HEARTBEAT_REQUEST_FROM_CONFIG_NODE, 
req.heartbeatId);
 
-    collectPipeMetaReport(logger, false).setTo(resp);
+    collectPipeMetaReport(logger).setTo(resp);
     PipeInsertionDataNodeListener.getInstance().listenToHeartbeat(true);
   }
 
-  private PipeMetaReport collectPipeMetaReport(
-      final Optional<Logger> logger, final boolean includeQueryMode) throws 
TException {
-    final Set<Integer> dataRegionIds =
-        StorageEngine.getInstance().getAllDataRegionIds().stream()
-            .map(DataRegionId::getId)
-            .collect(Collectors.toSet());
-
+  private PipeMetaReport collectPipeMetaReport(final Optional<Logger> logger) 
throws TException {
     final PipeMetaReport report = new PipeMetaReport();
     try {
       for (final PipeMeta pipeMeta : pipeMetaKeeper.getPipeMetaList()) {
@@ -542,17 +537,23 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
         final PipeStaticMeta staticMeta = pipeMeta.getStaticMeta();
 
         final Map<Integer, PipeTask> pipeTaskMap = 
pipeTaskManager.getPipeTasks(staticMeta);
-        final boolean isAllDataRegionCompleted =
-            pipeTaskMap == null
-                || pipeTaskMap.entrySet().stream()
-                    .filter(entry -> dataRegionIds.contains(entry.getKey()))
-                    .allMatch(entry -> ((PipeDataNodeTask) 
entry.getValue()).isCompleted());
-        final boolean isCompleted =
-            isAllDataRegionCompleted && includeDataAndNeedDrop(pipeMeta, 
includeQueryMode);
+        final Set<Integer> expectedDataRegionIds = 
getExpectedDataRegionIds(pipeMeta);
+        final List<Integer> completedDataRegionIds = new ArrayList<>();
+        if (pipeTaskMap != null) {
+          for (final Integer regionId : expectedDataRegionIds) {
+            final PipeTask pipeTask = pipeTaskMap.get(regionId);
+            if (pipeTask instanceof PipeDataNodeTask
+                && ((PipeDataNodeTask) pipeTask).isCompleted()) {
+              completedDataRegionIds.add(regionId);
+            }
+          }
+        }
+        report.pipeCompletedDataRegionList.add(
+            new TPipeCompletedDataRegion(
+                staticMeta.getPipeName(), staticMeta.getCreationTime(), 
completedDataRegionIds));
         final Pair<Long, Double> remainingEventAndTime =
             PipeDataNodeSinglePipeMetrics.getInstance()
                 .getRemainingEventAndTime(staticMeta.getPipeName(), 
staticMeta.getCreationTime());
-        report.pipeCompletedList.add(isCompleted);
         
report.pipeRemainingEventCountList.add(remainingEventAndTime.getLeft());
         report.pipeRemainingTimeList.add(remainingEventAndTime.getRight());
         report.pipeDegradedStatusList.add(
@@ -561,16 +562,6 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
                     .getGlobalTsFileEpochDegraded()));
         report.pipeRecentFailureList.add(
             ((PipeTemporaryMetaInAgent) 
pipeMeta.getTemporaryMeta()).getRecentFailures());
-
-        logger.ifPresent(
-            l ->
-                PipeLogger.log(
-                    l::info,
-                    DataNodePipeMessages
-                        
.LOG_REPORTING_PIPE_META_ARG_ISCOMPLETED_ARG_REMAININGEVENTCOUNT_ARG_8F996DF3,
-                    pipeMeta.coreReportMessage(),
-                    isCompleted,
-                    remainingEventAndTime.getLeft()));
       }
       logger.ifPresent(
           l ->
@@ -578,56 +569,68 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
                   l::info,
                   DataNodePipeMessages.LOG_REPORTED_ARG_PIPE_METAS_12068FC6,
                   report.pipeMetaBinaryList.size()));
-    } catch (final IOException | IllegalPathException e) {
+    } catch (final IOException e) {
       throw new TException(e);
     }
     return report;
   }
 
-  private boolean includeDataAndNeedDrop(final PipeMeta pipeMeta, final 
boolean includeQueryMode)
-      throws IllegalPathException {
-    final PipeParameters sourceParameters = 
pipeMeta.getStaticMeta().getSourceParameters();
-    if 
(!DataRegionListeningFilter.parseInsertionDeletionListeningOptionPair(sourceParameters)
-        .getLeft()) {
-      return false;
-    }
-    if (!includeQueryMode) {
-      return isSnapshotMode(sourceParameters);
+  // Returns the DataRegion ids that this DataNode is expected to transfer for 
the given pipe.
+  // A region is included only when it is owned by this DataNode, is led by 
this DataNode according
+  // to the pipe's runtime metadata, and is selected by the pipe's source 
parameters. This expected
+  // set is used instead of the already-created PipeTask map so that a failed 
task initialization is
+  // not silently treated as a completed region.
+  private Set<Integer> getExpectedDataRegionIds(final PipeMeta pipeMeta) {
+    final PipeStaticMeta staticMeta = pipeMeta.getStaticMeta();
+    final PipeParameters sourceParameters = staticMeta.getSourceParameters();
+    final Set<Integer> localDataRegionIds =
+        StorageEngine.getInstance().getAllDataRegionIds().stream()
+            .map(DataRegionId::getId)
+            .collect(Collectors.toSet());
+    final Set<Integer> expectedDataRegionIds = new HashSet<>();
+    for (final Map.Entry<Integer, PipeTaskMeta> entry :
+        
pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().entrySet()) {
+      final int regionId = entry.getKey();
+      if (entry.getValue().getLeaderNodeId() != CONFIG.getDataNodeId()
+          || !localDataRegionIds.contains(regionId)) {
+        continue;
+      }
+      try {
+        if (DataRegionListeningFilter.shouldDataRegionBeListened(
+            sourceParameters, new DataRegionId(regionId), 
staticMeta.getPipeType())) {
+          expectedDataRegionIds.add(regionId);
+        }
+      } catch (final IllegalPathException e) {
+        throw new PipeException(e.toString());
+      }
     }
-
-    final String sourceModeValue =
-        sourceParameters.getStringOrDefault(
-            Arrays.asList(
-                PipeSourceConstant.EXTRACTOR_MODE_KEY, 
PipeSourceConstant.SOURCE_MODE_KEY),
-            PipeSourceConstant.EXTRACTOR_MODE_DEFAULT_VALUE);
-    return 
sourceModeValue.equalsIgnoreCase(PipeSourceConstant.EXTRACTOR_MODE_QUERY_VALUE)
-        || 
sourceModeValue.equalsIgnoreCase(PipeSourceConstant.EXTRACTOR_MODE_SNAPSHOT_VALUE);
+    return expectedDataRegionIds;
   }
 
   private static class PipeMetaReport {
     private final List<ByteBuffer> pipeMetaBinaryList = new ArrayList<>();
-    private final List<Boolean> pipeCompletedList = new ArrayList<>();
     private final List<Long> pipeRemainingEventCountList = new ArrayList<>();
     private final List<Double> pipeRemainingTimeList = new ArrayList<>();
     private final List<Integer> pipeDegradedStatusList = new ArrayList<>();
     private final List<Map<String, Long>> pipeRecentFailureList = new 
ArrayList<>();
+    private final List<TPipeCompletedDataRegion> pipeCompletedDataRegionList = 
new ArrayList<>();
 
     private void setTo(final TDataNodeHeartbeatResp resp) {
       resp.setPipeMetaList(pipeMetaBinaryList);
-      resp.setPipeCompletedList(pipeCompletedList);
       resp.setPipeRemainingEventCountList(pipeRemainingEventCountList);
       resp.setPipeRemainingTimeList(pipeRemainingTimeList);
       resp.setPipeDegradedStatusList(pipeDegradedStatusList);
       resp.setPipeRecentFailureList(pipeRecentFailureList);
+      resp.setPipeCompletedDataRegionList(pipeCompletedDataRegionList);
     }
 
     private void setTo(final TPipeHeartbeatResp resp) {
       resp.setPipeMetaList(pipeMetaBinaryList);
-      resp.setPipeCompletedList(pipeCompletedList);
       resp.setPipeRemainingEventCountList(pipeRemainingEventCountList);
       resp.setPipeRemainingTimeList(pipeRemainingTimeList);
       resp.setPipeDegradedStatusList(pipeDegradedStatusList);
       resp.setPipeRecentFailureList(pipeRecentFailureList);
+      resp.setPipeCompletedDataRegionList(pipeCompletedDataRegionList);
     }
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/util/ModsOperationUtil.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/util/ModsOperationUtil.java
index e2b65b5415c..38487b693e9 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/util/ModsOperationUtil.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/util/ModsOperationUtil.java
@@ -19,8 +19,10 @@
 
 package org.apache.iotdb.db.pipe.event.common.tsfile.parser.util;
 
+import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.path.PatternTreeMap;
 import org.apache.iotdb.db.i18n.DataNodePipeMessages;
+import 
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionPathUtils;
 import org.apache.iotdb.db.storageengine.dataregion.modification.ModEntry;
 import 
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile;
 import org.apache.iotdb.db.utils.ModificationUtils;
@@ -90,7 +92,7 @@ public class ModsOperationUtil {
       return false;
     }
 
-    final List<ModEntry> mods = modifications.getOverlapped(deviceID, 
measurementID);
+    final List<ModEntry> mods = getOverlappedMods(deviceID, measurementID, 
modifications);
     if (mods == null || mods.isEmpty()) {
       return false;
     }
@@ -127,7 +129,7 @@ public class ModsOperationUtil {
     List<ModsInfo> modsInfos = new ArrayList<>(measurements.size());
 
     for (final String measurement : measurements) {
-      final List<ModEntry> mods = modifications.getOverlapped(deviceID, 
measurement);
+      final List<ModEntry> mods = getOverlappedMods(deviceID, measurement, 
modifications);
       if (mods == null || mods.isEmpty()) {
         // No mods, use empty list and index 0
         modsInfos.add(new ModsInfo(Collections.emptyList(), 0));
@@ -156,6 +158,17 @@ public class ModsOperationUtil {
     return modsInfos;
   }
 
+  private static List<ModEntry> getOverlappedMods(
+      final IDeviceID deviceID,
+      final String measurement,
+      final PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> 
modifications) {
+    try {
+      return modifications.getOverlapped(CompactionPathUtils.getPath(deviceID, 
measurement));
+    } catch (final IllegalPathException e) {
+      throw new PipeException(e.getMessage(), e);
+    }
+  }
+
   /**
    * Check if data at the specified time point is deleted
    *
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTemporaryMetaInCoordinator.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTemporaryMetaInCoordinator.java
index 5c3a1ea17eb..79fc8b6a5b2 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTemporaryMetaInCoordinator.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTemporaryMetaInCoordinator.java
@@ -33,7 +33,7 @@ import java.util.concurrent.ConcurrentMap;
 public class PipeTemporaryMetaInCoordinator implements PipeTemporaryMeta {
 
   // ConfigNode statistics
-  private final Set<Integer> completedDataNodeIds =
+  private final Set<Integer> completedDataRegionIds =
       Collections.newSetFromMap(new ConcurrentHashMap<>());
   private final ConcurrentMap<Integer, Long> nodeId2RemainingEventMap = new 
ConcurrentHashMap<>();
   private final ConcurrentMap<Integer, Double> nodeId2RemainingTimeMap = new 
ConcurrentHashMap<>();
@@ -41,12 +41,12 @@ public class PipeTemporaryMetaInCoordinator implements 
PipeTemporaryMeta {
   private final ConcurrentMap<Integer, RecentFailureSnapshot> 
nodeId2RecentFailuresMap =
       new ConcurrentHashMap<>();
 
-  public void markDataNodeCompleted(final int dataNodeId) {
-    completedDataNodeIds.add(dataNodeId);
+  public void markDataRegionCompleted(final int dataRegionId) {
+    completedDataRegionIds.add(dataRegionId);
   }
 
-  public void markDataNodeUncompleted(final int dataNodeId) {
-    completedDataNodeIds.remove(dataNodeId);
+  public Set<Integer> getCompletedDataRegionIds() {
+    return completedDataRegionIds;
   }
 
   public void setRemainingEvent(final int dataNodeId, final long 
remainingEventCount) {
@@ -86,10 +86,6 @@ public class PipeTemporaryMetaInCoordinator implements 
PipeTemporaryMeta {
     }
   }
 
-  public Set<Integer> getCompletedDataNodeIds() {
-    return completedDataNodeIds;
-  }
-
   public long getGlobalRemainingEvents() {
     return 
nodeId2RemainingEventMap.values().stream().reduce(Long::sum).orElse(0L);
   }
@@ -131,7 +127,7 @@ public class PipeTemporaryMetaInCoordinator implements 
PipeTemporaryMeta {
       return false;
     }
     final PipeTemporaryMetaInCoordinator that = 
(PipeTemporaryMetaInCoordinator) o;
-    return Objects.equals(this.completedDataNodeIds, that.completedDataNodeIds)
+    return Objects.equals(this.completedDataRegionIds, 
that.completedDataRegionIds)
         && Objects.equals(this.nodeId2RemainingEventMap, 
that.nodeId2RemainingEventMap)
         && Objects.equals(this.nodeId2RemainingTimeMap, 
that.nodeId2RemainingTimeMap)
         && Objects.equals(this.nodeId2IsDegradedMap, that.nodeId2IsDegradedMap)
@@ -141,7 +137,7 @@ public class PipeTemporaryMetaInCoordinator implements 
PipeTemporaryMeta {
   @Override
   public int hashCode() {
     return Objects.hash(
-        completedDataNodeIds,
+        completedDataRegionIds,
         nodeId2RemainingEventMap,
         nodeId2RemainingTimeMap,
         nodeId2IsDegradedMap,
@@ -151,8 +147,8 @@ public class PipeTemporaryMetaInCoordinator implements 
PipeTemporaryMeta {
   @Override
   public String toString() {
     return "PipeTemporaryMeta{"
-        + "completedDataNodeIds="
-        + completedDataNodeIds
+        + "completedDataRegionIds="
+        + completedDataRegionIds
         + ", nodeId2RemainingEventMap="
         + nodeId2RemainingEventMap
         + ", nodeId2RemainingTimeMap="
diff --git a/iotdb-protocol/thrift-commons/src/main/thrift/common.thrift 
b/iotdb-protocol/thrift-commons/src/main/thrift/common.thrift
index edc824a6b43..16b40a85354 100644
--- a/iotdb-protocol/thrift-commons/src/main/thrift/common.thrift
+++ b/iotdb-protocol/thrift-commons/src/main/thrift/common.thrift
@@ -197,6 +197,12 @@ struct TSetThrottleQuotaReq {
   2: required TThrottleQuota throttleQuota
 }
 
+struct TPipeCompletedDataRegion {
+  1: required string pipeName
+  2: required i64 creationTime
+  3: required list<i32> completedDataRegionIds
+}
+
 struct TPipeHeartbeatResp {
   1: required list<binary> pipeMetaList
   2: optional list<bool> pipeCompletedList
@@ -204,6 +210,7 @@ struct TPipeHeartbeatResp {
   4: optional list<double> pipeRemainingTimeList
   5: optional list<i32> pipeDegradedStatusList
   6: optional list<map<string, i64>> pipeRecentFailureList
+  7: optional list<TPipeCompletedDataRegion> pipeCompletedDataRegionList
 }
 
 struct TLicense {
diff --git a/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift 
b/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift
index d66ed10ccf9..1c3fd78470d 100644
--- a/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift
+++ b/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift
@@ -319,6 +319,7 @@ struct TDataNodeHeartbeatResp {
   17: optional map<i32, i64> dataRegionRawDataSize
   18: optional list<i32> pipeDegradedStatusList
   19: optional list<map<string, i64>> pipeRecentFailureList
+  20: optional list<common.TPipeCompletedDataRegion> 
pipeCompletedDataRegionList
 }
 
 struct TPipeHeartbeatReq {

Reply via email to