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

jackietien pushed a commit to branch rc/1.3.5
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 02e358b31343ff0fc1cf598ed3290b8b640ff2e0
Author: Zhenyu Luo <[email protected]>
AuthorDate: Thu Aug 7 14:41:31 2025 +0800

    [To dev/1.3] Pipe: Delete the heartbeat event count in Remaining Count 
#16115 (#16116)
    
    * Pipe: Delete the heartbeat event count in Remaining Count
    
    * update
    
    * delete getRemainingEvents function
---
 .../PipeDataNodeRemainingEventAndTimeOperator.java      | 17 ++---------------
 .../metric/overview/PipeDataNodeSinglePipeMetrics.java  |  4 ++--
 2 files changed, 4 insertions(+), 17 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java
index 0308e9b5b63..099b84d99dc 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java
@@ -116,21 +116,8 @@ public class PipeDataNodeRemainingEventAndTimeOperator 
extends PipeRemainingOper
     return remainingEvents >= 0 ? remainingEvents : 0;
   }
 
-  long getRemainingEvents() {
-    final long remainingEvents =
-        tsfileEventCount.get()
-            + rawTabletEventCount.get()
-            + insertNodeEventCount.get()
-            + heartbeatEventCount.get()
-            + schemaRegionExtractors.stream()
-                .map(IoTDBSchemaRegionExtractor::getUnTransferredEventCount)
-                .reduce(Long::sum)
-                .orElse(0L);
-
-    // There are cases where the indicator is negative. For example, after the 
Pipe is restarted,
-    // the Processor SubTask is still collecting Events, resulting in a 
negative count. This
-    // situation cannot be avoided because the Pipe may be restarted 
internally.
-    return remainingEvents >= 0 ? remainingEvents : 0;
+  public int getInsertNodeEventCount() {
+    return insertNodeEventCount.get();
   }
 
   /**
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java
index 1840093a347..1cbbae7ec37 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java
@@ -89,7 +89,7 @@ public class PipeDataNodeSinglePipeMetrics implements 
IMetricSet {
         Metric.PIPE_DATANODE_REMAINING_EVENT_COUNT.toString(),
         MetricLevel.IMPORTANT,
         operator,
-        PipeDataNodeRemainingEventAndTimeOperator::getRemainingEvents,
+        
PipeDataNodeRemainingEventAndTimeOperator::getRemainingNonHeartbeatEvents,
         Tag.NAME.toString(),
         operator.getPipeName(),
         Tag.CREATION_TIME.toString(),
@@ -399,7 +399,7 @@ public class PipeDataNodeSinglePipeMetrics implements 
IMetricSet {
         remainingEventAndTimeOperatorMap.computeIfAbsent(
             pipeName + "_" + creationTime,
             k -> new PipeDataNodeRemainingEventAndTimeOperator(pipeName, 
creationTime));
-    return new Pair<>(operator.getRemainingEvents(), 
operator.getRemainingTime());
+    return new Pair<>(operator.getRemainingNonHeartbeatEvents(), 
operator.getRemainingTime());
   }
 
   //////////////////////////// singleton ////////////////////////////

Reply via email to