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

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


The following commit(s) were added to refs/heads/rc/1.3.3 by this push:
     new cb0d2a51650 Pipe: fix calculations of 
PipeDataNodeRemainingEventAndTimeOperator (#13876) (#13889)
cb0d2a51650 is described below

commit cb0d2a51650e1e32cb878962a428a347d90dc1e5
Author: V_Galaxy <[email protected]>
AuthorDate: Wed Oct 23 19:50:20 2024 +0800

    Pipe: fix calculations of PipeDataNodeRemainingEventAndTimeOperator 
(#13876) (#13889)
---
 .../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 bd01f26b655..5932c6b544b 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
@@ -136,7 +136,7 @@ public class PipeInsertNodeTabletInsertionEvent extends 
EnrichedEvent
       PipeDataNodeResourceManager.wal().pin(walEntryHandler);
       if (Objects.nonNull(pipeName)) {
         PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
-            .increaseInsertionEventCount(pipeName + "_" + creationTime);
+            .increaseTabletEventCount(pipeName + "_" + creationTime);
       }
       return true;
     } catch (final Exception e) {
@@ -169,7 +169,7 @@ public class PipeInsertNodeTabletInsertionEvent extends 
EnrichedEvent
     } 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 a5795ede49c..b89c11b1f39 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
@@ -118,7 +118,7 @@ public class PipeRawTabletInsertionEvent extends 
EnrichedEvent implements Tablet
                 PipeMemoryWeightUtil.calculateTabletSizeInBytes(tablet));
     if (Objects.nonNull(pipeName)) {
       PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
-          .increaseInsertionEventCount(pipeName + "_" + creationTime);
+          .increaseTabletEventCount(pipeName + "_" + creationTime);
     }
     return true;
   }
@@ -127,7 +127,7 @@ public class PipeRawTabletInsertionEvent extends 
EnrichedEvent implements Tablet
   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