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

Reply via email to