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 d7888e84f28 Pipe: eliminate progress index for raw tablet event to
avoid endless rebind (#13542)
d7888e84f28 is described below
commit d7888e84f289072059d55b47a851b2028985c32e
Author: V_Galaxy <[email protected]>
AuthorDate: Fri Sep 20 17:45:52 2024 +0800
Pipe: eliminate progress index for raw tablet event to avoid endless rebind
(#13542)
---
.../event/common/tablet/PipeRawTabletInsertionEvent.java | 6 +++++-
.../pipe/event/common/tsfile/PipeTsFileInsertionEvent.java | 6 +++++-
.../realtime/assigner/PipeDataRegionAssigner.java | 13 +++++++------
.../assigner/PipeTimePartitionProgressIndexKeeper.java | 6 ++++--
4 files changed, 21 insertions(+), 10 deletions(-)
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 f06819323ec..64f53e1cd1b 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
@@ -25,6 +25,7 @@ import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
import org.apache.iotdb.commons.pipe.pattern.PipePattern;
import org.apache.iotdb.commons.pipe.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.utils.TestOnly;
+import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryWeightUtil;
import org.apache.iotdb.db.pipe.resource.memory.PipeTabletMemoryBlock;
@@ -51,7 +52,7 @@ public class PipeRawTabletInsertionEvent extends
EnrichedEvent implements Tablet
private TabletInsertionDataContainer dataContainer;
- private ProgressIndex overridingProgressIndex;
+ private volatile ProgressIndex overridingProgressIndex;
private PipeRawTabletInsertionEvent(
final Tablet tablet,
@@ -136,6 +137,9 @@ public class PipeRawTabletInsertionEvent extends
EnrichedEvent implements Tablet
protected void reportProgress() {
if (needToReport) {
super.reportProgress();
+ if (sourceEvent instanceof PipeTsFileInsertionEvent) {
+ ((PipeTsFileInsertionEvent) sourceEvent).eliminateProgressIndex();
+ }
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index 9d7e80f6fdc..3bf863d0442 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -72,7 +72,7 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
implements TsFileIns
// May be updated after it is flushed. Should be negative if not set.
private long flushPointCount = TsFileProcessor.FLUSH_POINT_COUNT_NOT_SET;
- private ProgressIndex overridingProgressIndex;
+ private volatile ProgressIndex overridingProgressIndex;
public PipeTsFileInsertionEvent(
final TsFileResource resource,
@@ -302,6 +302,10 @@ public class PipeTsFileInsertionEvent extends
EnrichedEvent implements TsFileIns
@Override
protected void reportProgress() {
super.reportProgress();
+ this.eliminateProgressIndex();
+ }
+
+ public void eliminateProgressIndex() {
if (Objects.isNull(overridingProgressIndex)) {
PipeTimePartitionProgressIndexKeeper.getInstance()
.eliminateProgressIndex(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeDataRegionAssigner.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeDataRegionAssigner.java
index 849c2eecf56..88f1d53b9ea 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeDataRegionAssigner.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeDataRegionAssigner.java
@@ -38,6 +38,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.Closeable;
+import java.util.concurrent.atomic.AtomicReference;
public class PipeDataRegionAssigner implements Closeable {
@@ -59,7 +60,8 @@ public class PipeDataRegionAssigner implements Closeable {
private int counter = 0;
- private ProgressIndex maxProgressIndexForTsFileInsertionEvent =
MinimumProgressIndex.INSTANCE;
+ private final AtomicReference<ProgressIndex>
maxProgressIndexForTsFileInsertionEvent =
+ new AtomicReference<>(MinimumProgressIndex.INSTANCE);
public String getDataRegionId() {
return dataRegionId;
@@ -160,16 +162,15 @@ public class PipeDataRegionAssigner implements Closeable {
if (PipeTimePartitionProgressIndexKeeper.getInstance()
.isProgressIndexAfterOrEquals(
dataRegionId, event.getTimePartitionId(),
event.getProgressIndex())) {
-
event.bindProgressIndex(maxProgressIndexForTsFileInsertionEvent.deepCopy());
+
event.bindProgressIndex(maxProgressIndexForTsFileInsertionEvent.get().deepCopy());
LOGGER.warn(
"Data region {} bind {} to event {} because it was flushed
prematurely.",
dataRegionId,
maxProgressIndexForTsFileInsertionEvent,
- event);
+ event.coreReportMessage());
} else {
- maxProgressIndexForTsFileInsertionEvent =
-
maxProgressIndexForTsFileInsertionEvent.updateToMinimumEqualOrIsAfterProgressIndex(
- event.getProgressIndex());
+ maxProgressIndexForTsFileInsertionEvent.updateAndGet(
+ index ->
index.updateToMinimumEqualOrIsAfterProgressIndex(event.getProgressIndex()));
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeTimePartitionProgressIndexKeeper.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeTimePartitionProgressIndexKeeper.java
index f4f81e48fe1..df563eba6bb 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeTimePartitionProgressIndexKeeper.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/assigner/PipeTimePartitionProgressIndexKeeper.java
@@ -25,6 +25,7 @@ import org.apache.tsfile.utils.Pair;
import java.util.Map;
import java.util.Map.Entry;
+import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
public class PipeTimePartitionProgressIndexKeeper {
@@ -58,7 +59,7 @@ public class PipeTimePartitionProgressIndexKeeper {
if (v == null) {
return null;
}
- if (v.getRight() && v.getLeft().equals(progressIndex)) {
+ if (v.getRight() && !v.getLeft().isAfter(progressIndex)) {
return new Pair<>(v.getLeft(), false);
}
return v;
@@ -75,7 +76,8 @@ public class PipeTimePartitionProgressIndexKeeper {
.map(Entry::getValue)
.filter(pair -> pair.right)
.map(Pair::getLeft)
- .anyMatch(index -> progressIndex.isAfter(index) ||
progressIndex.equals(index));
+ .filter(Objects::nonNull)
+ .anyMatch(index -> !index.isAfter(progressIndex));
}
//////////////////////////// singleton ////////////////////////////