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

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


The following commit(s) were added to refs/heads/dev/1.3 by this push:
     new 2b7f32de1eb Pipe: Reduced the log of epoch switching & Refactor & 
Optimized the memory reservation logic of tsFile parser provider & Subscription 
IT: assertGte for received tsfile count (#15068) (#15070)
2b7f32de1eb is described below

commit 2b7f32de1eb66bac19a9d377548ce8291622d573
Author: Caideyipi <[email protected]>
AuthorDate: Mon Mar 31 18:46:15 2025 +0800

    Pipe: Reduced the log of epoch switching & Refactor & Optimized the memory 
reservation logic of tsFile parser provider & Subscription IT: assertGte for 
received tsfile count (#15068) (#15070)
    
    Co-authored-by: VGalaxies <[email protected]>
---
 .../IoTDBPathLooseDeviceTsfilePushConsumerIT.java  | 10 +++++-----
 .../IoTDBTimeLooseTsfilePushConsumerIT.java        | 10 +++++-----
 .../IoTDBTSPatternTsfilePushConsumerIT.java        |  2 +-
 .../time/IoTDBRealTimeDBTsfilePushConsumerIT.java  |  8 ++++----
 .../time/IoTDBTimeRangeDBTsfilePushConsumerIT.java | 22 ++++++++++------------
 .../async/IoTDBDataRegionAsyncConnector.java       | 21 ---------------------
 .../common/tablet/PipeRawTabletInsertionEvent.java |  4 ++--
 .../TsFileInsertionDataContainerProvider.java      | 19 ++++++++++++-------
 .../PipeRealtimeDataRegionHybridExtractor.java     | 12 ++++++++----
 .../pipe/resource/tsfile/PipeTsFileResource.java   |  2 +-
 10 files changed, 48 insertions(+), 62 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java
index f6a65093dbd..d1a65470935 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java
@@ -198,7 +198,7 @@ public class IoTDBPathLooseDeviceTsfilePushConsumerIT 
extends AbstractSubscripti
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1);
+          assertGte(onReceive.get(), 1);
           assertGte(rowCounts.get(0).get(), 3, "Write data before 
subscription" + device);
           assertGte(rowCounts.get(0).get(), 3, "Write data before 
subscription" + device2);
         });
@@ -210,7 +210,7 @@ public class IoTDBPathLooseDeviceTsfilePushConsumerIT 
extends AbstractSubscripti
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1);
+          assertGte(onReceive.get(), 1);
           assertGte(rowCounts.get(0).get(), 3, "Write out-of-range data" + 
device);
           assertGte(rowCounts.get(0).get(), 3, "Write out-of-range data" + 
device2);
         });
@@ -222,7 +222,7 @@ public class IoTDBPathLooseDeviceTsfilePushConsumerIT 
extends AbstractSubscripti
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 2);
+          assertGte(onReceive.get(), 2);
           assertGte(rowCounts.get(0).get(), 8, "write data" + device);
           assertGte(rowCounts.get(0).get(), 8, "write data " + device2);
         });
@@ -234,7 +234,7 @@ public class IoTDBPathLooseDeviceTsfilePushConsumerIT 
extends AbstractSubscripti
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 3);
+          assertGte(onReceive.get(), 3);
           assertGte(rowCounts.get(0).get(), 10, "Write data: end boundary at " 
+ device);
           assertGte(rowCounts.get(0).get(), 10, "Write data: end boundary at " 
+ device2);
         });
@@ -245,7 +245,7 @@ public class IoTDBPathLooseDeviceTsfilePushConsumerIT 
extends AbstractSubscripti
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 3);
+          assertGte(onReceive.get(), 3);
           assertGte(rowCounts.get(0).get(), 10, "Write data: > end " + device);
           assertGte(rowCounts.get(0).get(), 10, "Write data: > end " + 
device2);
         });
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java
index 79286c6c23c..50f37c4183e 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java
@@ -185,7 +185,7 @@ public class IoTDBTimeLooseTsfilePushConsumerIT extends 
AbstractSubscriptionRegr
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1);
+          assertGte(onReceive.get(), 1);
           assertGte(rowCount.get(), 3);
         });
 
@@ -196,7 +196,7 @@ public class IoTDBTimeLooseTsfilePushConsumerIT extends 
AbstractSubscriptionRegr
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1);
+          assertGte(onReceive.get(), 1);
           assertGte(rowCount.get(), 3);
         });
 
@@ -207,7 +207,7 @@ public class IoTDBTimeLooseTsfilePushConsumerIT extends 
AbstractSubscriptionRegr
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 2);
+          assertGte(onReceive.get(), 2);
           assertGte(rowCount.get(), 8);
         });
 
@@ -218,7 +218,7 @@ public class IoTDBTimeLooseTsfilePushConsumerIT extends 
AbstractSubscriptionRegr
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 3);
+          assertGte(onReceive.get(), 3);
           assertGte(rowCount.get(), 10);
         });
 
@@ -229,7 +229,7 @@ public class IoTDBTimeLooseTsfilePushConsumerIT extends 
AbstractSubscriptionRegr
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 3);
+          assertGte(onReceive.get(), 3);
           assertGte(rowCount.get(), 10);
         });
   }
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/pattern/IoTDBTSPatternTsfilePushConsumerIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/pattern/IoTDBTSPatternTsfilePushConsumerIT.java
index 67ec1b32852..f12000d89b6 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/pattern/IoTDBTSPatternTsfilePushConsumerIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/pattern/IoTDBTSPatternTsfilePushConsumerIT.java
@@ -217,7 +217,7 @@ public class IoTDBTSPatternTsfilePushConsumerIT extends 
AbstractSubscriptionRegr
     System.out.println(FORMAT.format(new Date()) + " src:" + 
getCount(session_src, sql));
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceiveCount.get(), 2, "receive files over 2");
+          assertGte(onReceiveCount.get(), 2, "receive files over 2");
           assertEquals(rowCounts.get(0).get(), 25, device + ".s_0");
           assertEquals(rowCounts.get(1).get(), 0, device + ".s_1");
           assertEquals(rowCounts.get(2).get(), 0, database + ".d_1.s_0");
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBRealTimeDBTsfilePushConsumerIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBRealTimeDBTsfilePushConsumerIT.java
index 4af7d1acd4a..c6b9420f13b 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBRealTimeDBTsfilePushConsumerIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBRealTimeDBTsfilePushConsumerIT.java
@@ -163,8 +163,8 @@ public class IoTDBRealTimeDBTsfilePushConsumerIT extends 
AbstractSubscriptionReg
     insert_data(System.currentTimeMillis());
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1, "should process 1 file");
-          assertEquals(rowCount.get(), 4, "4 records");
+          assertGte(onReceive.get(), 1, "should process 1 file");
+          assertGte(rowCount.get(), 4, "4 records");
         });
 
     // Subscribe and then write data
@@ -172,8 +172,8 @@ public class IoTDBRealTimeDBTsfilePushConsumerIT extends 
AbstractSubscriptionReg
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 2, "should process 2 file");
-          assertEquals(rowCount.get(), 8, "8 records");
+          assertGte(onReceive.get(), 2, "should process 2 file");
+          assertGte(rowCount.get(), 8, "8 records");
         });
   }
 }
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBTimeRangeDBTsfilePushConsumerIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBTimeRangeDBTsfilePushConsumerIT.java
index 448bb26821a..2a77ca0256a 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBTimeRangeDBTsfilePushConsumerIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/regression/pushconsumer/time/IoTDBTimeRangeDBTsfilePushConsumerIT.java
@@ -161,38 +161,36 @@ public class IoTDBTimeRangeDBTsfilePushConsumerIT extends 
AbstractSubscriptionRe
 
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1);
-          // loose-time should 2 records,get 4 records
-          assertTrue(rowCount.get() >= 2);
+          assertGte(onReceive.get(), 1);
+          assertGte(rowCount.get(), 2);
         });
 
     insert_data(System.currentTimeMillis()); // now, not in range
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 1);
-          assertTrue(rowCount.get() >= 2);
+          assertGte(onReceive.get(), 1);
+          assertGte(rowCount.get(), 2);
         });
 
     insert_data(1707782400000L); // 2024-02-13 08:00:00+08:00
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 2);
-          assertTrue(rowCount.get() >= 6);
+          assertGte(onReceive.get(), 2);
+          assertGte(rowCount.get(), 6);
         });
 
     insert_data(1711814398000L); // 2024-03-30 23:59:58+08:00
     AWAIT.untilAsserted(
         () -> {
-          // Because the end time is 2024-03-31 00:00:00, closed interval
-          assertEquals(onReceive.get(), 3);
-          assertTrue(rowCount.get() >= 8);
+          assertGte(onReceive.get(), 3);
+          assertGte(rowCount.get(), 8);
         });
 
     insert_data(1711900798000L); // 2024-03-31 23:59:58+08:00
     AWAIT.untilAsserted(
         () -> {
-          assertEquals(onReceive.get(), 3);
-          assertTrue(rowCount.get() >= 8);
+          assertGte(onReceive.get(), 3);
+          assertGte(rowCount.get(), 8);
         });
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBDataRegionAsyncConnector.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBDataRegionAsyncConnector.java
index ca38765820c..4e71c3af431 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBDataRegionAsyncConnector.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBDataRegionAsyncConnector.java
@@ -677,27 +677,6 @@ public class IoTDBDataRegionAsyncConnector extends 
IoTDBConnector {
     return retryEventQueue.size();
   }
 
-  // For performance, this will not acquire lock and does not guarantee the 
correct
-  // result. However, this shall not cause any exceptions when concurrently 
read & written.
-  public int getRetryEventCount(final String pipeName) {
-    final AtomicInteger count = new AtomicInteger(0);
-    try {
-      retryEventQueue.forEach(
-          event -> {
-            if (event instanceof EnrichedEvent
-                && pipeName.equals(((EnrichedEvent) event).getPipeName())) {
-              count.incrementAndGet();
-            }
-          });
-      return count.get();
-    } catch (final Exception e) {
-      if (LOGGER.isDebugEnabled()) {
-        LOGGER.debug("Failed to get retry event count for pipe {}.", pipeName, 
e);
-      }
-      return count.get();
-    }
-  }
-
   //////////////////////// APIs provided for PipeTransferTrackableHandler 
////////////////////////
 
   public boolean isClosed() {
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 de4a36e1c65..f7884076f26 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
@@ -282,8 +282,8 @@ public class PipeRawTabletInsertionEvent extends 
EnrichedEvent
   }
 
   public long count() {
-    final Tablet covertedTablet = shouldParseTimeOrPattern() ? 
convertToTablet() : tablet;
-    return (long) covertedTablet.rowSize * covertedTablet.getSchemas().size();
+    final Tablet convertedTablet = shouldParseTimeOrPattern() ? 
convertToTablet() : tablet;
+    return (long) convertedTablet.rowSize * 
convertedTablet.getSchemas().size();
   }
 
   /////////////////////////// parsePatternOrTime ///////////////////////////
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/TsFileInsertionDataContainerProvider.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/TsFileInsertionDataContainerProvider.java
index cffe6e87d9c..5c8b87f9b20 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/TsFileInsertionDataContainerProvider.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/TsFileInsertionDataContainerProvider.java
@@ -27,6 +27,7 @@ import 
org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
 import 
org.apache.iotdb.db.pipe.event.common.tsfile.container.query.TsFileInsertionQueryDataContainer;
 import 
org.apache.iotdb.db.pipe.event.common.tsfile.container.scan.TsFileInsertionScanDataContainer;
 import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
+import org.apache.iotdb.db.pipe.resource.tsfile.PipeTsFileResource;
 
 import org.apache.tsfile.file.metadata.IDeviceID;
 import org.apache.tsfile.file.metadata.PlainDeviceID;
@@ -63,13 +64,17 @@ public class TsFileInsertionDataContainerProvider {
   }
 
   public TsFileInsertionDataContainer provide() throws IOException {
-    if (startTime != Long.MIN_VALUE
-        || endTime != Long.MAX_VALUE
-        || pattern instanceof IoTDBPipePattern
-            && !((IoTDBPipePattern) 
pattern).mayMatchMultipleTimeSeriesInOneDevice()) {
-      // 1. If time filter exists, use query here because the scan container 
may filter it
-      // row by row in single page chunk.
-      // 2. If the pattern matches only one time series in one device, use 
query container here
+    // Use scan container to save memory
+    if ((double) 
PipeDataNodeResourceManager.memory().getUsedMemorySizeInBytes()
+            / PipeDataNodeResourceManager.memory().getTotalMemorySizeInBytes()
+        > PipeTsFileResource.MEMORY_SUFFICIENT_THRESHOLD) {
+      return new TsFileInsertionScanDataContainer(
+          tsFile, pattern, startTime, endTime, pipeTaskMeta, sourceEvent);
+    }
+
+    if (pattern instanceof IoTDBPipePattern
+        && !((IoTDBPipePattern) 
pattern).mayMatchMultipleTimeSeriesInOneDevice()) {
+      // If the pattern matches only one time series in one device, use query 
container here
       // because there is no timestamps merge overhead.
       //
       // Note: We judge prefix pattern as matching multiple timeseries in one 
device because it's
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/PipeRealtimeDataRegionHybridExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/PipeRealtimeDataRegionHybridExtractor.java
index b1baca7c42a..cecdf23642d 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/PipeRealtimeDataRegionHybridExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/realtime/PipeRealtimeDataRegionHybridExtractor.java
@@ -82,13 +82,17 @@ public class PipeRealtimeDataRegionHybridExtractor extends 
PipeRealtimeDataRegio
   }
 
   private void extractTabletInsertion(final PipeRealtimeEvent event) {
-    if (canNotUseTabletAnyMore(event)) {
+    TsFileEpoch.State state = event.getTsFileEpoch().getState(this);
+
+    if (state != TsFileEpoch.State.USING_TSFILE
+        && state != TsFileEpoch.State.USING_BOTH
+        && canNotUseTabletAnyMore(event)) {
       event
           .getTsFileEpoch()
           .migrateState(
               this,
-              state -> {
-                switch (state) {
+              curState -> {
+                switch (curState) {
                   case EMPTY:
                   case USING_TSFILE:
                     return TsFileEpoch.State.USING_TSFILE;
@@ -100,7 +104,7 @@ public class PipeRealtimeDataRegionHybridExtractor extends 
PipeRealtimeDataRegio
               });
     }
 
-    final TsFileEpoch.State state = event.getTsFileEpoch().getState(this);
+    state = event.getTsFileEpoch().getState(this);
     switch (state) {
       case USING_TSFILE:
         // Ignore the tablet event.
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java
index 2cb3a59c399..dd3ca9c5bc0 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java
@@ -49,7 +49,7 @@ public class PipeTsFileResource implements AutoCloseable {
   private static final Logger LOGGER = 
LoggerFactory.getLogger(PipeTsFileResource.class);
 
   public static final long TSFILE_MIN_TIME_TO_LIVE_IN_MS = 1000L * 20;
-  private static final float MEMORY_SUFFICIENT_THRESHOLD = 0.7f;
+  public static final float MEMORY_SUFFICIENT_THRESHOLD = 0.7f;
 
   private final File hardlinkOrCopiedFile;
   private final boolean isTsFile;

Reply via email to