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 -> {