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;