From 3771af76041313c15f45efe65d3ffded101684b2 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Mon, 21 Sep 2026 10:45:36 +0800 Subject: [PATCH] [Pipe] Avoid no-op consensus writes for covered progress (#18574) (cherry picked from commit 5aae6697fba555c40a3bba40529ef8e7f8a2e2c2) --- .../heartbeat/PipeHeartbeatParser.java | 20 +++++++------------ 1 file changed, 7 insertions(+), 13 deletions(-) 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 ace07f5e2d3d3..bfdb951b4093e 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 @@ -192,18 +192,12 @@ private void parseHeartbeatAndSaveMetaChangeLocally( } // Update progress index - if (!(runtimeMetaFromCoordinator - .getValue() - .getProgressIndex() - .isAfter(runtimeMetaFromAgent.getProgressIndex()) - || runtimeMetaFromCoordinator - .getValue() - .getProgressIndex() - .equals(runtimeMetaFromAgent.getProgressIndex()))) { + final ProgressIndex coordinatorProgressIndex = + runtimeMetaFromCoordinator.getValue().getProgressIndex(); + final ProgressIndex agentProgressIndex = runtimeMetaFromAgent.getProgressIndex(); + if (!coordinatorProgressIndex.isEqualOrAfter(agentProgressIndex)) { final ProgressIndex updatedProgressIndex = - runtimeMetaFromCoordinator - .getValue() - .updateProgressIndex(runtimeMetaFromAgent.getProgressIndex()); + runtimeMetaFromCoordinator.getValue().updateProgressIndex(agentProgressIndex); PipeConfigNodeResourceManager.log() .schedule( PipeHeartbeatParser.class, @@ -217,8 +211,8 @@ private void parseHeartbeatAndSaveMetaChangeLocally( + "Progress index on coordinator: {}, progress index from agent: {}, updated progressIndex: {}", pipeMetaFromCoordinator.getStaticMeta().getPipeName(), runtimeMetaFromCoordinator.getKey(), - runtimeMetaFromCoordinator.getValue().getProgressIndex(), - runtimeMetaFromAgent.getProgressIndex(), + coordinatorProgressIndex, + agentProgressIndex, updatedProgressIndex)); needWriteConsensusOnConfigNodes.set(true);