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 {