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

rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 0c7a9471f0e Pipe: fix calculations of 
PipeDataNodeRemainingEventAndTimeOperator (#13876)
0c7a9471f0e is described below

commit 0c7a9471f0ea84a0f6435a7a5009a7a439bd2786
Author: V_Galaxy <[email protected]>
AuthorDate: Wed Oct 23 18:21:43 2024 +0800

    Pipe: fix calculations of PipeDataNodeRemainingEventAndTimeOperator (#13876)
---
 .../common/tablet/PipeInsertNodeTabletInsertionEvent.java   |  4 ++--
 .../event/common/tablet/PipeRawTabletInsertionEvent.java    |  4 ++--
 .../metric/PipeDataNodeRemainingEventAndTimeMetrics.java    |  8 ++++----
 .../metric/PipeDataNodeRemainingEventAndTimeOperator.java   | 13 +++++++------
 4 files changed, 15 insertions(+), 14 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
index 7fb8a0b470e..88c91016f6f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
@@ -163,7 +163,7 @@ public class PipeInsertNodeTabletInsertionEvent extends 
PipeInsertionEvent
       PipeDataNodeResourceManager.wal().pin(walEntryHandler);
       if (Objects.nonNull(pipeName)) {
         PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
-            .increaseInsertionEventCount(pipeName + "_" + creationTime);
+            .increaseTabletEventCount(pipeName + "_" + creationTime);
       }
       return true;
     } catch (final Exception e) {
@@ -196,7 +196,7 @@ public class PipeInsertNodeTabletInsertionEvent extends 
PipeInsertionEvent
     } finally {
       if (Objects.nonNull(pipeName)) {
         PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
-            .decreaseInsertionEventCount(pipeName + "_" + creationTime);
+            .decreaseTabletEventCount(pipeName + "_" + creationTime);
       }
     }
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
index 3fccf9ba516..e205738753e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
@@ -171,7 +171,7 @@ public class PipeRawTabletInsertionEvent extends 
PipeInsertionEvent
                 PipeMemoryWeightUtil.calculateTabletSizeInBytes(tablet));
     if (Objects.nonNull(pipeName)) {
       PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
-          .increaseInsertionEventCount(pipeName + "_" + creationTime);
+          .increaseTabletEventCount(pipeName + "_" + creationTime);
     }
     return true;
   }
@@ -180,7 +180,7 @@ public class PipeRawTabletInsertionEvent extends 
PipeInsertionEvent
   public boolean internallyDecreaseResourceReferenceCount(final String 
holderMessage) {
     if (Objects.nonNull(pipeName)) {
       PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
-          .decreaseInsertionEventCount(pipeName + "_" + creationTime);
+          .decreaseTabletEventCount(pipeName + "_" + creationTime);
     }
     allocatedMemoryBlock.close();
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeMetrics.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeMetrics.java
index a8765b3b61d..25c2ede2407 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeMetrics.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeMetrics.java
@@ -129,16 +129,16 @@ public class PipeDataNodeRemainingEventAndTimeMetrics 
implements IMetricSet {
     }
   }
 
-  public void increaseInsertionEventCount(final String pipeID) {
+  public void increaseTabletEventCount(final String pipeID) {
     remainingEventAndTimeOperatorMap
         .computeIfAbsent(pipeID, k -> new 
PipeDataNodeRemainingEventAndTimeOperator())
-        .increaseInsertionEventCount();
+        .increaseTabletEventCount();
   }
 
-  public void decreaseInsertionEventCount(final String pipeID) {
+  public void decreaseTabletEventCount(final String pipeID) {
     remainingEventAndTimeOperatorMap
         .computeIfAbsent(pipeID, k -> new 
PipeDataNodeRemainingEventAndTimeOperator())
-        .decreaseInsertionEventCount();
+        .decreaseTabletEventCount();
   }
 
   public void increaseTsFileEventCount(final String pipeID) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeOperator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeOperator.java
index c0072215146..bee0e6975b4 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeOperator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/PipeDataNodeRemainingEventAndTimeOperator.java
@@ -44,7 +44,7 @@ class PipeDataNodeRemainingEventAndTimeOperator extends 
PipeRemainingOperator {
   private final Set<IoTDBSchemaRegionExtractor> schemaRegionExtractors =
       Collections.newSetFromMap(new ConcurrentHashMap<>());
 
-  private final AtomicInteger insertionEventCount = new AtomicInteger(0);
+  private final AtomicInteger tabletEventCount = new AtomicInteger(0);
   private final AtomicInteger tsfileEventCount = new AtomicInteger(0);
   private final AtomicInteger heartbeatEventCount = new AtomicInteger(0);
 
@@ -58,12 +58,12 @@ class PipeDataNodeRemainingEventAndTimeOperator extends 
PipeRemainingOperator {
 
   //////////////////////////// Remaining event & time calculation 
////////////////////////////
 
-  void increaseInsertionEventCount() {
-    tsfileEventCount.incrementAndGet();
+  void increaseTabletEventCount() {
+    tabletEventCount.incrementAndGet();
   }
 
-  void decreaseInsertionEventCount() {
-    tsfileEventCount.decrementAndGet();
+  void decreaseTabletEventCount() {
+    tabletEventCount.decrementAndGet();
   }
 
   void increaseTsFileEventCount() {
@@ -84,6 +84,7 @@ class PipeDataNodeRemainingEventAndTimeOperator extends 
PipeRemainingOperator {
 
   long getRemainingEvents() {
     return tsfileEventCount.get()
+        + tabletEventCount.get()
         + heartbeatEventCount.get()
         + schemaRegionExtractors.stream()
             .map(IoTDBSchemaRegionExtractor::getUnTransferredEventCount)
@@ -105,7 +106,7 @@ class PipeDataNodeRemainingEventAndTimeOperator extends 
PipeRemainingOperator {
     final double invocationValue = collectInvocationHistogram.getMean();
     // Do not take heartbeat event into account
     final double totalDataRegionWriteEventCount =
-        tsfileEventCount.get() * Math.max(invocationValue, 1) + 
insertionEventCount.get();
+        tsfileEventCount.get() * Math.max(invocationValue, 1) + 
tabletEventCount.get();
 
     dataRegionCommitMeter.updateAndGet(
         meter -> {

Reply via email to