This is an automated email from the ASF dual-hosted git repository.
jt2594838 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 6ca8850a25b [Pipe] Support configurable TsFile parser (#18440)
6ca8850a25b is described below
commit 6ca8850a25bb333b5345d094ceea14511aab5a61
Author: Caideyipi <[email protected]>
AuthorDate: Wed Aug 12 09:33:35 2026 +0800
[Pipe] Support configurable TsFile parser (#18440)
---
.../tsfile/PipeCompactedTsFileInsertionEvent.java | 1 +
.../common/tsfile/PipeTsFileInsertionEvent.java | 61 ++++++++++++---------
.../parser/TsFileInsertionEventParserProvider.java | 62 ++++++++++++++++++++++
.../source/dataregion/IoTDBDataRegionSource.java | 17 ++++++
...istoricalDataRegionTsFileAndDeletionSource.java | 6 +++
.../realtime/PipeRealtimeDataRegionSource.java | 9 ++++
.../realtime/assigner/PipeDataRegionAssigner.java | 1 +
.../pipe/event/TsFileInsertionEventParserTest.java | 41 ++++++++++++++
.../db/pipe/source/IoTDBDataRegionSourceTest.java | 31 +++++++++++
.../pipe/config/constant/PipeSourceConstant.java | 5 ++
10 files changed, 210 insertions(+), 24 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeCompactedTsFileInsertionEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeCompactedTsFileInsertionEvent.java
index a0f1d52b967..83a37b261f0 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeCompactedTsFileInsertionEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeCompactedTsFileInsertionEvent.java
@@ -93,6 +93,7 @@ public class PipeCompactedTsFileInsertionEvent extends
PipeTsFileInsertionEvent
flushPointCount = bindFlushPointCount(originalEvents);
overridingProgressIndex = bindOverridingProgressIndex(originalEvents);
bindTsFileDedupScopeID(anyOfOriginalEvents.getTsFileDedupScopeID());
+ setTsFileParser(anyOfOriginalEvents.getTsFileParser());
}
private static boolean bindIsWithMod(Set<PipeTsFileInsertionEvent>
originalEvents) {
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 b969ac4e618..335278e2f37 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
@@ -119,6 +119,7 @@ public class PipeTsFileInsertionEvent extends
PipeInsertionEvent
protected volatile ProgressIndex overridingProgressIndex;
private Set<String> tableNames;
private String tsFileDedupScopeID;
+ private String tsFileParser;
// False when generated tablet events should wait for an external progress
report.
private volatile boolean shouldReportGeneratedEventsOnCommit = true;
@@ -540,6 +541,14 @@ public class PipeTsFileInsertionEvent extends
PipeInsertionEvent
return tsFileDedupScopeID;
}
+ public String getTsFileParser() {
+ return tsFileParser;
+ }
+
+ public void setTsFileParser(final String tsFileParser) {
+ this.tsFileParser = tsFileParser;
+ }
+
@Override
public PipeTsFileInsertionEvent
shallowCopySelfAndBindPipeTaskMetaForProgressReport(
final String pipeName,
@@ -553,29 +562,32 @@ public class PipeTsFileInsertionEvent extends
PipeInsertionEvent
final boolean skipIfNoPrivileges,
final long startTime,
final long endTime) {
- return new PipeTsFileInsertionEvent(
- getRawIsTableModelEvent(),
- getSourceDatabaseNameFromDataRegion(),
- resource,
- tsFile,
- isWithMod,
- isLoaded,
- isGeneratedByHistoricalExtractor,
- tableNames,
- pipeName,
- creationTime,
- pipeTaskMeta,
- treePattern,
- tablePattern,
- userId,
- userName,
- cliHostname,
- skipIfNoPrivileges,
- startTime,
- endTime,
- isTsFileSealed)
- .bindTsFileDedupScopeID(tsFileDedupScopeID)
-
.setShouldReportGeneratedEventsOnCommit(shouldReportGeneratedEventsOnCommit);
+ final PipeTsFileInsertionEvent copiedEvent =
+ new PipeTsFileInsertionEvent(
+ getRawIsTableModelEvent(),
+ getSourceDatabaseNameFromDataRegion(),
+ resource,
+ tsFile,
+ isWithMod,
+ isLoaded,
+ isGeneratedByHistoricalExtractor,
+ tableNames,
+ pipeName,
+ creationTime,
+ pipeTaskMeta,
+ treePattern,
+ tablePattern,
+ userId,
+ userName,
+ cliHostname,
+ skipIfNoPrivileges,
+ startTime,
+ endTime,
+ isTsFileSealed)
+ .bindTsFileDedupScopeID(tsFileDedupScopeID)
+
.setShouldReportGeneratedEventsOnCommit(shouldReportGeneratedEventsOnCommit);
+ copiedEvent.setTsFileParser(tsFileParser);
+ return copiedEvent;
}
@Override
@@ -1140,7 +1152,8 @@ public class PipeTsFileInsertionEvent extends
PipeInsertionEvent
shouldParse4Privilege
? new UserEntity(Long.parseLong(userId), userName,
cliHostname)
: null,
- this)
+ this,
+ tsFileParser)
.provide(isWithMod));
return eventParser.get();
} catch (final Exception e) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParserProvider.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParserProvider.java
index c08c68307a2..0a3205c112e 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParserProvider.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParserProvider.java
@@ -42,6 +42,9 @@ import java.util.Map;
import java.util.Objects;
import java.util.stream.Collectors;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_QUERY_VALUE;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_SCAN_VALUE;
+
public class TsFileInsertionEventParserProvider {
private final String pipeName;
@@ -56,6 +59,7 @@ public class TsFileInsertionEventParserProvider {
protected final PipeTaskMeta pipeTaskMeta;
protected final PipeTsFileInsertionEvent sourceEvent;
private final IAuditEntity entity;
+ private final String tsFileParser;
public TsFileInsertionEventParserProvider(
final String pipeName,
@@ -68,6 +72,32 @@ public class TsFileInsertionEventParserProvider {
final PipeTaskMeta pipeTaskMeta,
final IAuditEntity entity,
final PipeTsFileInsertionEvent sourceEvent) {
+ this(
+ pipeName,
+ creationTime,
+ tsFile,
+ treePattern,
+ tablePattern,
+ startTime,
+ endTime,
+ pipeTaskMeta,
+ entity,
+ sourceEvent,
+ null);
+ }
+
+ public TsFileInsertionEventParserProvider(
+ final String pipeName,
+ final long creationTime,
+ final File tsFile,
+ final TreePattern treePattern,
+ final TablePattern tablePattern,
+ final long startTime,
+ final long endTime,
+ final PipeTaskMeta pipeTaskMeta,
+ final IAuditEntity entity,
+ final PipeTsFileInsertionEvent sourceEvent,
+ final String tsFileParser) {
this.pipeName = pipeName;
this.creationTime = creationTime;
this.tsFile = tsFile;
@@ -78,6 +108,7 @@ public class TsFileInsertionEventParserProvider {
this.pipeTaskMeta = pipeTaskMeta;
this.entity = entity;
this.sourceEvent = sourceEvent;
+ this.tsFileParser = tsFileParser;
}
public TsFileInsertionEventParser provide(final boolean isWithMod)
@@ -101,6 +132,37 @@ public class TsFileInsertionEventParserProvider {
isWithMod);
}
+ if (EXTRACTOR_TSFILE_PARSER_QUERY_VALUE.equals(tsFileParser)) {
+ return new TsFileInsertionEventQueryParser(
+ pipeName,
+ creationTime,
+ tsFile,
+ treePattern,
+ startTime,
+ endTime,
+ pipeTaskMeta,
+ sourceEvent,
+ entity,
+ sourceEvent.isSkipIfNoPrivileges(),
+ null,
+ isWithMod);
+ }
+
+ if (EXTRACTOR_TSFILE_PARSER_SCAN_VALUE.equals(tsFileParser)) {
+ return new TsFileInsertionEventScanParser(
+ pipeName,
+ creationTime,
+ tsFile,
+ treePattern,
+ startTime,
+ endTime,
+ pipeTaskMeta,
+ entity,
+ sourceEvent.isSkipIfNoPrivileges(),
+ sourceEvent,
+ isWithMod);
+ }
+
// Use scan container to save memory
if ((double)
PipeDataNodeResourceManager.memory().getUsedMemorySizeInBytes()
/
PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes()
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java
index 6adcb1addb5..5d677e10047 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java
@@ -97,6 +97,9 @@ import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.E
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_START_TIME_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TABLE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TABLE_NAME_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_QUERY_VALUE;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_SCAN_VALUE;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_WATERMARK_INTERVAL_DEFAULT_VALUE;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_WATERMARK_INTERVAL_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_DATABASE_KEY;
@@ -120,6 +123,7 @@ import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.S
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_START_TIME_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_TABLE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_TABLE_NAME_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_TSFILE_PARSER_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_WATERMARK_INTERVAL_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant._EXTRACTOR_WATERMARK_INTERVAL_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant._SOURCE_WATERMARK_INTERVAL_KEY;
@@ -165,6 +169,19 @@ public class IoTDBDataRegionSource extends IoTDBSource {
EXTRACTOR_PATTERN_FORMAT_PREFIX_VALUE,
EXTRACTOR_PATTERN_FORMAT_IOTDB_VALUE);
+ // If unset, the parser is selected automatically.
+ validator
+ .validateAttributeValueRange(
+ EXTRACTOR_TSFILE_PARSER_KEY,
+ true,
+ EXTRACTOR_TSFILE_PARSER_QUERY_VALUE,
+ EXTRACTOR_TSFILE_PARSER_SCAN_VALUE)
+ .validateAttributeValueRange(
+ SOURCE_TSFILE_PARSER_KEY,
+ true,
+ EXTRACTOR_TSFILE_PARSER_QUERY_VALUE,
+ EXTRACTOR_TSFILE_PARSER_SCAN_VALUE);
+
// Validate tree pattern and table pattern
validatePattern(TreePattern.parsePipePatternFromSourceParameters(validator.getParameters()));
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
index cf913176df8..1ee89144d88 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
@@ -111,6 +111,7 @@ import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.E
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_MODS_ENABLE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_MODS_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_START_TIME_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_END_TIME_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_HISTORY_ENABLE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_HISTORY_END_TIME_KEY;
@@ -121,6 +122,7 @@ import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.S
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_MODS_ENABLE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_MODS_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_START_TIME_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_TSFILE_PARSER_KEY;
import static
org.apache.iotdb.commons.pipe.source.IoTDBSource.getSkipIfNoPrivileges;
import static org.apache.tsfile.common.constant.TsFileConstant.PATH_ROOT;
import static org.apache.tsfile.common.constant.TsFileConstant.PATH_SEPARATOR;
@@ -163,6 +165,7 @@ public class PipeHistoricalDataRegionTsFileAndDeletionSource
private boolean shouldExtractInsertion;
private boolean shouldExtractDeletion;
private boolean shouldTransferModFile; // Whether to transfer mods
+ private String tsFileParser;
protected String userId;
protected String userName;
protected String cliHostname;
@@ -366,6 +369,8 @@ public class PipeHistoricalDataRegionTsFileAndDeletionSource
treePattern = TreePattern.parsePipePatternFromSourceParameters(parameters);
tablePattern =
TablePattern.parsePipePatternFromSourceParameters(parameters);
+ tsFileParser =
+ parameters.getStringByKeys(EXTRACTOR_TSFILE_PARSER_KEY,
SOURCE_TSFILE_PARSER_KEY);
final DataRegion dataRegion =
StorageEngine.getInstance().getDataRegion(new
DataRegionId(environment.getRegionId()));
@@ -1285,6 +1290,7 @@ public class
PipeHistoricalDataRegionTsFileAndDeletionSource
skipIfNoPrivileges,
historicalDataExtractionStartTime,
historicalDataExtractionEndTime);
+ event.setTsFileParser(tsFileParser);
if (shouldUseHistoricalTsFileQueryPriorityOrder()) {
event.skipReportOnCommitAndGeneratedEvents();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java
index 732275824eb..e4da82c465c 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java
@@ -84,11 +84,13 @@ import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.E
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_REALTIME_LOOSE_RANGE_PATH_VALUE;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_REALTIME_LOOSE_RANGE_TIME_VALUE;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_START_TIME_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_TSFILE_PARSER_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_END_TIME_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_MODS_ENABLE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_MODS_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_REALTIME_LOOSE_RANGE_KEY;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_START_TIME_KEY;
+import static
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_TSFILE_PARSER_KEY;
import static
org.apache.iotdb.commons.pipe.source.IoTDBSource.getSkipIfNoPrivileges;
public abstract class PipeRealtimeDataRegionSource implements PipeExtractor {
@@ -123,6 +125,7 @@ public abstract class PipeRealtimeDataRegionSource
implements PipeExtractor {
protected boolean isForwardingPipeRequests;
private boolean shouldTransferModFile; // Whether to transfer mods
+ private String tsFileParser;
private boolean sloppyTimeRange; // true to disable time range filter after
extraction
private boolean sloppyPattern; // true to disable pattern filter after
extraction
@@ -239,6 +242,8 @@ public abstract class PipeRealtimeDataRegionSource
implements PipeExtractor {
treePattern = TreePattern.parsePipePatternFromSourceParameters(parameters);
tablePattern =
TablePattern.parsePipePatternFromSourceParameters(parameters);
+ tsFileParser =
+ parameters.getStringByKeys(EXTRACTOR_TSFILE_PARSER_KEY,
SOURCE_TSFILE_PARSER_KEY);
final DataRegion dataRegion =
StorageEngine.getInstance().getDataRegion(new
DataRegionId(environment.getRegionId()));
@@ -615,6 +620,10 @@ public abstract class PipeRealtimeDataRegionSource
implements PipeExtractor {
return shouldTransferModFile;
}
+ public final String getTsFileParser() {
+ return tsFileParser;
+ }
+
private void maySkipProgressIndexForRealtimeEvent(final PipeRealtimeEvent
event) {
if (PipeTsFileEpochProgressIndexKeeper.getInstance()
.isProgressIndexAfterOrEquals(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java
index 6e40be3d324..4c717da3c53 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java
@@ -192,6 +192,7 @@ public class PipeDataRegionAssigner implements Closeable {
final PipeTsFileInsertionEvent tsFileInsertionEvent =
(PipeTsFileInsertionEvent) innerEvent;
tsFileInsertionEvent.bindTsFileDedupScopeID(source.getTsFileDedupScopeID());
+ tsFileInsertionEvent.setTsFileParser(source.getTsFileParser());
tsFileInsertionEvent.disableMod4NonTransferPipes(source.isShouldTransferModFile());
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionEventParserTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionEventParserTest.java
index f12664ff8b5..0a4e0421abe 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionEventParserTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionEventParserTest.java
@@ -31,6 +31,7 @@ import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.TsFileInsertionEventParser;
import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.TsFileInsertionEventParserMemoryBlock;
+import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.TsFileInsertionEventParserProvider;
import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.query.TsFileInsertionEventQueryParser;
import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.scan.AlignedSinglePageWholeChunkReader;
import
org.apache.iotdb.db.pipe.event.common.tsfile.parser.scan.SinglePageWholeChunkReader;
@@ -157,6 +158,46 @@ public class TsFileInsertionEventParserTest {
System.out.println(System.currentTimeMillis() - startTime);
}
+ @Test
+ public void testConfiguredTsFileParserSelection() throws Exception {
+ final PipeTsFileInsertionEvent sourceEvent =
+
createPipeTsFileInsertionEventForRetryTest("configured-parser-selection.tsfile");
+
+ try (final TsFileInsertionEventParser parser =
+ new TsFileInsertionEventParserProvider(
+ null,
+ 0,
+ nonalignedTsFile,
+ new PrefixTreePattern("root"),
+ null,
+ Long.MIN_VALUE,
+ Long.MAX_VALUE,
+ null,
+ null,
+ sourceEvent,
+ "query")
+ .provide(false)) {
+ Assert.assertTrue(parser instanceof TsFileInsertionEventQueryParser);
+ }
+
+ try (final TsFileInsertionEventParser parser =
+ new TsFileInsertionEventParserProvider(
+ null,
+ 0,
+ nonalignedTsFile,
+ new PrefixTreePattern("root"),
+ null,
+ Long.MIN_VALUE,
+ Long.MAX_VALUE,
+ null,
+ null,
+ sourceEvent,
+ "scan")
+ .provide(false)) {
+ Assert.assertTrue(parser instanceof TsFileInsertionEventScanParser);
+ }
+ }
+
@Test
public void testScanParserReleasesTabletMemoryAfterRawTabletGenerated()
throws Exception {
nonalignedTsFile =
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/IoTDBDataRegionSourceTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/IoTDBDataRegionSourceTest.java
index dd0558bc238..f660423b0b0 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/IoTDBDataRegionSourceTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/IoTDBDataRegionSourceTest.java
@@ -23,6 +23,7 @@ import
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
import org.apache.iotdb.db.pipe.source.dataregion.IoTDBDataRegionSource;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
+import org.apache.iotdb.pipe.api.exception.PipeParameterNotValidException;
import org.junit.Assert;
import org.junit.Test;
@@ -52,6 +53,36 @@ public class IoTDBDataRegionSourceTest {
}
}
+ @Test
+ public void testTsFileParserParameter() throws Exception {
+ for (final String parser : new String[] {"query", "scan"}) {
+ try (final IoTDBDataRegionSource extractor = new
IoTDBDataRegionSource()) {
+ extractor.validate(
+ new PipeParameterValidator(
+ new PipeParameters(
+ new HashMap<String, String>() {
+ {
+ put(PipeSourceConstant.SOURCE_TSFILE_PARSER_KEY,
parser);
+ }
+ })));
+ }
+ }
+
+ try (final IoTDBDataRegionSource extractor = new IoTDBDataRegionSource()) {
+ Assert.assertThrows(
+ PipeParameterNotValidException.class,
+ () ->
+ extractor.validate(
+ new PipeParameterValidator(
+ new PipeParameters(
+ new HashMap<String, String>() {
+ {
+ put(PipeSourceConstant.SOURCE_TSFILE_PARSER_KEY,
"invalid");
+ }
+ }))));
+ }
+ }
+
@Test
public void testIoTDBDataRegionExtractorWithPattern() {
Assert.assertEquals(
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
index 9e0f5a2ecce..055553c7e1a 100644
---
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
@@ -120,6 +120,11 @@ public class PipeSourceConstant {
public static final String SOURCE_MODS_KEY = "source.mods";
public static final boolean EXTRACTOR_MODS_DEFAULT_VALUE =
EXTRACTOR_MODS_ENABLE_DEFAULT_VALUE;
+ public static final String EXTRACTOR_TSFILE_PARSER_KEY =
"extractor.tsfile.parser";
+ public static final String SOURCE_TSFILE_PARSER_KEY = "source.tsfile.parser";
+ public static final String EXTRACTOR_TSFILE_PARSER_QUERY_VALUE = "query";
+ public static final String EXTRACTOR_TSFILE_PARSER_SCAN_VALUE = "scan";
+
public static final String EXTRACTOR_REALTIME_ENABLE_KEY =
"extractor.realtime.enable";
public static final String SOURCE_REALTIME_ENABLE_KEY =
"source.realtime.enable";
public static final boolean EXTRACTOR_REALTIME_ENABLE_DEFAULT_VALUE = true;