From a7f39581a8103fdc95dbf0fd7b78ddbaf45c0bdf Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Mon, 8 Jun 2026 14:57:44 +0800 Subject: [PATCH] Improve object column parser compatibility --- .../parser/TsFileInsertionEventParser.java | 39 ++++++++ .../TsFileInsertionEventQueryParser.java | 94 +++++++++++++++++-- .../scan/TsFileInsertionEventScanParser.java | 58 +++++++++++- .../TsFileInsertionEventTableParser.java | 60 +++++++++++- ...sertionEventTableParserTabletIterator.java | 88 ++++++++++++++++- .../handler/PipeTransferTsFileHandler.java | 25 +++++ .../plan/node/load/LoadSingleTsFileNode.java | 45 ++++++++- .../planner/plan/node/write/ObjectNode.java | 13 +++ .../load/splitter/TsFileSplitter.java | 13 +++ 9 files changed, 425 insertions(+), 10 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParser.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParser.java index 627731fa7edc4..1011b5c05db88 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParser.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/TsFileInsertionEventParser.java @@ -33,6 +33,7 @@ import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryBlock; import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryWeightUtil; import org.apache.iotdb.db.storageengine.dataregion.modification.ModEntry; +import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; import org.apache.iotdb.db.utils.datastructure.PatternTreeMapFactory; import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent; @@ -45,6 +46,7 @@ import java.io.File; import java.io.IOException; +import java.util.List; public abstract class TsFileInsertionEventParser implements AutoCloseable { @@ -78,6 +80,8 @@ public abstract class TsFileInsertionEventParser implements AutoCloseable { protected Iterable tabletInsertionIterable; + protected final boolean objectPathsOnly; + protected TsFileInsertionEventParser( final File tsFile, final String pipeName, @@ -91,6 +95,38 @@ protected TsFileInsertionEventParser( final boolean skipIfNoPrivileges, final PipeInsertionEvent sourceEvent, final boolean isWithMod) { + this( + tsFile, + pipeName, + creationTime, + treePattern, + tablePattern, + startTime, + endTime, + pipeTaskMeta, + entity, + skipIfNoPrivileges, + sourceEvent, + null, + false, + isWithMod); + } + + protected TsFileInsertionEventParser( + final File tsFile, + final String pipeName, + final long creationTime, + final TreePattern treePattern, + final TablePattern tablePattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final IAuditEntity entity, + final boolean skipIfNoPrivileges, + final PipeInsertionEvent sourceEvent, + final TsFileResource tsFileResource, + final boolean objectPathsOnly, + final boolean isWithMod) { this.pipeName = pipeName; this.creationTime = creationTime; this.entity = entity; @@ -107,6 +143,7 @@ protected TsFileInsertionEventParser( this.pipeTaskMeta = pipeTaskMeta; this.sourceEvent = sourceEvent; + this.objectPathsOnly = objectPathsOnly; this.allocatedMemoryBlockForTablet = PipeDataNodeResourceManager.memory() @@ -130,6 +167,8 @@ protected TsFileInsertionEventParser( */ public abstract Iterable toTabletInsertionEvents(); + public void drainGeneratedObjectColumnModEntriesTo(final List modEntries) {} + /** * Record parse start time when hasNext() is called for the first time and returns true. Should be * called in Iterator.hasNext() when it's the first call. diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/query/TsFileInsertionEventQueryParser.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/query/TsFileInsertionEventQueryParser.java index 2656ec7d72d4b..e092ebe45ec01 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/query/TsFileInsertionEventQueryParser.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/query/TsFileInsertionEventQueryParser.java @@ -90,7 +90,7 @@ public TsFileInsertionEventQueryParser( final long endTime, final PipeInsertionEvent sourceEvent) throws IOException, IllegalPathException { - this(null, 0, tsFile, pattern, startTime, endTime, null, sourceEvent, false); + this(null, 0, tsFile, pattern, startTime, endTime, null, sourceEvent, false, false); } public TsFileInsertionEventQueryParser( @@ -104,6 +104,31 @@ public TsFileInsertionEventQueryParser( final PipeInsertionEvent sourceEvent, final boolean isWithMod) throws IOException, IllegalPathException { + this( + pipeName, + creationTime, + tsFile, + pattern, + startTime, + endTime, + pipeTaskMeta, + sourceEvent, + isWithMod, + false); + } + + public TsFileInsertionEventQueryParser( + final String pipeName, + final long creationTime, + final File tsFile, + final TreePattern pattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final PipeInsertionEvent sourceEvent, + final boolean isWithMod, + final boolean objectPathsOnly) + throws IOException, IllegalPathException { this( pipeName, creationTime, @@ -116,7 +141,8 @@ public TsFileInsertionEventQueryParser( null, false, null, - isWithMod); + isWithMod, + objectPathsOnly); } public TsFileInsertionEventQueryParser( @@ -133,6 +159,37 @@ public TsFileInsertionEventQueryParser( final Map deviceIsAlignedMap, final boolean isWithMod) throws IOException, IllegalPathException { + this( + pipeName, + creationTime, + tsFile, + pattern, + startTime, + endTime, + pipeTaskMeta, + sourceEvent, + entity, + skipIfNoPrivileges, + deviceIsAlignedMap, + isWithMod, + false); + } + + public TsFileInsertionEventQueryParser( + final String pipeName, + final long creationTime, + final File tsFile, + final TreePattern pattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final PipeInsertionEvent sourceEvent, + final IAuditEntity entity, + final boolean skipIfNoPrivileges, + final Map deviceIsAlignedMap, + final boolean isWithMod, + final boolean objectPathsOnly) + throws IOException, IllegalPathException { super( tsFile, pipeName, @@ -145,6 +202,8 @@ public TsFileInsertionEventQueryParser( entity, skipIfNoPrivileges, sourceEvent, + null, + objectPathsOnly, isWithMod); try { @@ -157,7 +216,7 @@ public TsFileInsertionEventQueryParser( .forceAllocateForTabletWithRetry(currentModifications.ramBytesUsed()); final PipeTsFileResourceManager tsFileResourceManager = PipeDataNodeResourceManager.tsfile(); - final Map> deviceMeasurementsMap; + Map> deviceMeasurementsMap; // TsFileReader is not thread-safe, so we need to create it here and close it later. long memoryRequiredInBytes = @@ -197,6 +256,10 @@ public TsFileInsertionEventQueryParser( memoryRequiredInBytes += PipeMemoryWeightUtil.memoryOfIDeviceID2StrList(deviceMeasurementsMap); } + if (objectPathsOnly) { + deviceMeasurementsMap = + filterObjectMeasurementsMap(deviceMeasurementsMap, measurementDataTypeMap); + } allocatedMemoryBlock = PipeDataNodeResourceManager.memory().forceAllocate(memoryRequiredInBytes); @@ -373,9 +436,10 @@ private Map> readFilteredDeviceMeasurementsMap( final Map> result = new HashMap<>(); for (final IDeviceID device : devices) { - tsFileSequenceReader - .readDeviceMetadata(device) - .values() + tsFileSequenceReader.readDeviceMetadata(device).values().stream() + .filter( + timeseriesMetadata -> + !objectPathsOnly || timeseriesMetadata.getTsDataType() == TSDataType.OBJECT) .forEach( timeseriesMetadata -> result @@ -386,6 +450,24 @@ private Map> readFilteredDeviceMeasurementsMap( return result; } + private Map> filterObjectMeasurementsMap( + final Map> deviceMeasurementsMap, + final Map measurementDataTypeMap) { + final Map> result = new HashMap<>(); + for (final Map.Entry> entry : deviceMeasurementsMap.entrySet()) { + final List measurements = new ArrayList<>(); + for (final String measurement : entry.getValue()) { + if (measurementDataTypeMap.get(entry.getKey() + "." + measurement) == TSDataType.OBJECT) { + measurements.add(measurement); + } + } + if (!measurements.isEmpty()) { + result.put(entry.getKey(), measurements); + } + } + return result; + } + @Override public Iterable toTabletInsertionEvents() { if (tabletInsertionIterable == null) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/scan/TsFileInsertionEventScanParser.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/scan/TsFileInsertionEventScanParser.java index d843ab38ab48a..6ac3470382849 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/scan/TsFileInsertionEventScanParser.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/scan/TsFileInsertionEventScanParser.java @@ -123,6 +123,35 @@ public TsFileInsertionEventScanParser( final PipeInsertionEvent sourceEvent, final boolean isWithMod) throws IOException, IllegalPathException { + this( + pipeName, + creationTime, + tsFile, + pattern, + startTime, + endTime, + pipeTaskMeta, + entity, + skipIfNoPrivileges, + sourceEvent, + isWithMod, + false); + } + + public TsFileInsertionEventScanParser( + final String pipeName, + final long creationTime, + final File tsFile, + final TreePattern pattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final IAuditEntity entity, + final boolean skipIfNoPrivileges, + final PipeInsertionEvent sourceEvent, + final boolean isWithMod, + final boolean objectPathsOnly) + throws IOException, IllegalPathException { super( tsFile, pipeName, @@ -135,6 +164,8 @@ public TsFileInsertionEventScanParser( entity, skipIfNoPrivileges, sourceEvent, + null, + objectPathsOnly, isWithMod); this.startTime = startTime; @@ -181,6 +212,19 @@ public TsFileInsertionEventScanParser( final PipeInsertionEvent sourceEvent, final boolean isWithMod) throws IOException, IllegalPathException { + this(tsFile, pattern, startTime, endTime, pipeTaskMeta, sourceEvent, isWithMod, false); + } + + public TsFileInsertionEventScanParser( + final File tsFile, + final TreePattern pattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final PipeInsertionEvent sourceEvent, + final boolean isWithMod, + final boolean objectPathsOnly) + throws IOException, IllegalPathException { this( null, 0, @@ -192,7 +236,8 @@ public TsFileInsertionEventScanParser( null, false, sourceEvent, - isWithMod); + isWithMod, + objectPathsOnly); } @Override @@ -416,8 +461,12 @@ private boolean putValueToColumns(final BatchData data, final Tablet tablet, fin case TEXT: case BLOB: case STRING: + case OBJECT: PipeTabletUtils.putValue( tablet, rowIndex, i, tablet.getSchemas().get(i).getType(), Binary.EMPTY_VALUE); + break; + default: + break; } PipeTabletUtils.markNullValue(tablet, rowIndex, i); continue; @@ -469,6 +518,7 @@ private boolean putValueToColumns(final BatchData data, final Tablet tablet, fin case TEXT: case BLOB: case STRING: + case OBJECT: final Binary binary = primitiveType.getBinary(); PipeTabletUtils.putValue( tablet, @@ -524,6 +574,7 @@ private boolean putValueToColumns(final BatchData data, final Tablet tablet, fin case TEXT: case BLOB: case STRING: + case OBJECT: final Binary binary = data.getBinary(); PipeTabletUtils.putValue( tablet, @@ -773,6 +824,11 @@ private boolean filterChunk( return true; } + if (objectPathsOnly && chunkHeader.getDataType() != TSDataType.OBJECT) { + tsFileSequenceReader.position(nextMarkerOffset); + return true; + } + // Skip the chunk if it is fully deleted by mods if (!currentModifications.isEmpty()) { final Statistics statistics = diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParser.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParser.java index 8ecdcc0cec5e9..d9f27c949516b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParser.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParser.java @@ -71,6 +71,35 @@ public TsFileInsertionEventTableParser( final PipeInsertionEvent sourceEvent, final boolean isWithMod) throws IOException { + this( + pipeName, + creationTime, + tsFile, + pattern, + startTime, + endTime, + pipeTaskMeta, + entity, + sourceEvent, + isWithMod, + false, + false); + } + + public TsFileInsertionEventTableParser( + final String pipeName, + final long creationTime, + final File tsFile, + final TablePattern pattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final IAuditEntity entity, + final PipeInsertionEvent sourceEvent, + final boolean isWithMod, + final boolean objectPathsOnly, + final boolean collectObjectColumnModEntries) + throws IOException { super( tsFile, pipeName, @@ -83,6 +112,8 @@ public TsFileInsertionEventTableParser( entity, true, sourceEvent, + null, + objectPathsOnly, isWithMod); this.isWithMod = isWithMod; @@ -140,6 +171,32 @@ public TsFileInsertionEventTableParser( null, 0, tsFile, pattern, startTime, endTime, pipeTaskMeta, entity, sourceEvent, isWithMod); } + public TsFileInsertionEventTableParser( + final File tsFile, + final TablePattern pattern, + final long startTime, + final long endTime, + final PipeTaskMeta pipeTaskMeta, + final IAuditEntity entity, + final PipeInsertionEvent sourceEvent, + final boolean isWithMod, + final boolean objectPathsOnly) + throws IOException { + this( + null, + 0, + tsFile, + pattern, + startTime, + endTime, + pipeTaskMeta, + entity, + sourceEvent, + isWithMod, + objectPathsOnly, + false); + } + @Override public Iterable toTabletInsertionEvents() { if (tabletInsertionIterable == null) { @@ -229,7 +286,8 @@ && hasTablePrivilege(entry.getKey()), allocatedMemoryBlockForTableSchemas, currentModifications, startTime, - endTime); + endTime, + objectPathsOnly); } while (tabletIterator.hasNext()) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java index 5b50eb166bebb..60bf57c5c5573 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java @@ -62,6 +62,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.function.BiConsumer; import java.util.function.Predicate; import java.util.stream.Collectors; @@ -84,6 +85,8 @@ public class TsFileInsertionEventTableParserTabletIterator implements Iterator modifications; @@ -128,10 +131,69 @@ public TsFileInsertionEventTableParserTabletIterator( final long startTime, final long endTime) throws IOException { + this( + tsFileSequenceReader, + predicate, + allocatedMemoryBlockForTablet, + allocatedMemoryBlockForBatchData, + allocatedMemoryBlockForChunk, + allocatedMemoryBlockForChunkMeta, + allocatedMemoryBlockForTableSchema, + modifications, + startTime, + endTime, + false); + } + + public TsFileInsertionEventTableParserTabletIterator( + final TsFileSequenceReader tsFileSequenceReader, + final Predicate> predicate, + final PipeMemoryBlock allocatedMemoryBlockForTablet, + final PipeMemoryBlock allocatedMemoryBlockForBatchData, + final PipeMemoryBlock allocatedMemoryBlockForChunk, + final PipeMemoryBlock allocatedMemoryBlockForChunkMeta, + final PipeMemoryBlock allocatedMemoryBlockForTableSchema, + final PatternTreeMap modifications, + final long startTime, + final long endTime, + final boolean objectPathsOnly) + throws IOException { + this( + tsFileSequenceReader, + predicate, + allocatedMemoryBlockForTablet, + allocatedMemoryBlockForBatchData, + allocatedMemoryBlockForChunk, + allocatedMemoryBlockForChunkMeta, + allocatedMemoryBlockForTableSchema, + modifications, + startTime, + endTime, + objectPathsOnly, + false, + null); + } + + public TsFileInsertionEventTableParserTabletIterator( + final TsFileSequenceReader tsFileSequenceReader, + final Predicate> predicate, + final PipeMemoryBlock allocatedMemoryBlockForTablet, + final PipeMemoryBlock allocatedMemoryBlockForBatchData, + final PipeMemoryBlock allocatedMemoryBlockForChunk, + final PipeMemoryBlock allocatedMemoryBlockForChunkMeta, + final PipeMemoryBlock allocatedMemoryBlockForTableSchema, + final PatternTreeMap modifications, + final long startTime, + final long endTime, + final boolean objectPathsOnly, + final boolean collectObjectColumnModEntries, + final BiConsumer> tableObjectMeasurementsSink) + throws IOException { this.startTime = startTime; this.endTime = endTime; this.modifications = modifications; + this.objectPathsOnly = objectPathsOnly; this.reader = tsFileSequenceReader; this.metadataQuerier = new MetadataQuerierByFileImpl(reader); @@ -224,6 +286,22 @@ public boolean hasNext() { continue; } + if (objectPathsOnly) { + final Iterator valueChunkMetadataIterator = + alignedChunkMetadata.getValueChunkMetadataList().iterator(); + while (valueChunkMetadataIterator.hasNext()) { + final IChunkMetadata valueChunkMetadata = valueChunkMetadataIterator.next(); + if (valueChunkMetadata == null + || valueChunkMetadata.getDataType() != TSDataType.OBJECT) { + valueChunkMetadataIterator.remove(); + } + } + if (alignedChunkMetadata.getValueChunkMetadataList().isEmpty()) { + chunkMetadataIterator.remove(); + continue; + } + } + size += PipeMemoryWeightUtil.calculateAlignedChunkMetaBytesUsed(alignedChunkMetadata); if (allocatedMemoryBlockForChunkMeta.getMemoryUsageInBytes() < size) { @@ -395,10 +473,13 @@ private void initChunkReader(final AbstractAlignedChunkMetadata alignedChunkMeta measurementList.subList(deviceIdSize, measurementList.size()).clear(); dataTypeList.subList(deviceIdSize, dataTypeList.size()).clear(); - boolean hasSelectedField = fieldSchemaList.isEmpty(); + boolean hasSelectedField = !objectPathsOnly && fieldSchemaList.isEmpty(); boolean hasSelectedNonNullChunk = false; for (; offset < fieldSchemaList.size(); ++offset) { final IMeasurementSchema schema = fieldSchemaList.get(offset); + if (objectPathsOnly && schema.getType() != TSDataType.OBJECT) { + continue; + } final String measurementName = internMeasurementName(schema); if (isFieldDeletedByMods( measurementName, @@ -499,7 +580,11 @@ private boolean fillMeasurementValueColumns( case TEXT: case BLOB: case STRING: + case OBJECT: PipeTabletUtils.putValue(tablet, rowIndex, i, dataTypeList.get(i), Binary.EMPTY_VALUE); + break; + default: + break; } PipeTabletUtils.markNullValue(tablet, rowIndex, i); continue; @@ -539,6 +624,7 @@ private boolean fillMeasurementValueColumns( case TEXT: case BLOB: case STRING: + case OBJECT: Binary binary = primitiveType.getBinary(); PipeTabletUtils.putValue( tablet, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java index 3fa5e557a208e..7d53d514b4cda 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java @@ -103,6 +103,31 @@ public PipeTransferTsFileHandler( final boolean transferMod, final String dataBaseName) throws InterruptedException { + this( + connector, + pipeName2WeightMap, + events, + eventsReferenceCount, + eventsHadBeenAddedToRetryQueue, + tsFile, + modFile, + null, + transferMod, + dataBaseName); + } + + public PipeTransferTsFileHandler( + final IoTDBDataRegionAsyncSink connector, + final Map, Double> pipeName2WeightMap, + final List events, + final AtomicInteger eventsReferenceCount, + final AtomicBoolean eventsHadBeenAddedToRetryQueue, + final File tsFile, + final File modFile, + final File objectDir, + final boolean transferMod, + final String dataBaseName) + throws InterruptedException { super(connector); this.pipeName2WeightMap = pipeName2WeightMap; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java index 098b9fde361d1..af3395b573c7a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java @@ -63,6 +63,8 @@ public class LoadSingleTsFileNode extends WritePlanNode { private final boolean deleteAfterLoad; private final long writePointCount; private boolean needDecodeTsFile; + private final boolean tsFileContainsObjectColumn; + private File objectFileSearchRoot; private TRegionReplicaSet localRegionReplicaSet; @@ -74,6 +76,28 @@ public LoadSingleTsFileNode( final boolean deleteAfterLoad, final long writePointCount, final boolean needDecodeTsFile) { + this( + id, + resource, + isTableModel, + database, + deleteAfterLoad, + writePointCount, + needDecodeTsFile, + false, + null); + } + + public LoadSingleTsFileNode( + final PlanNodeId id, + final TsFileResource resource, + final boolean isTableModel, + final String database, + final boolean deleteAfterLoad, + final long writePointCount, + final boolean needDecode4TimeColumn, + final boolean tsFileContainsObjectColumn, + final File objectFileSearchRoot) { super(id); this.tsFile = resource.getTsFile(); this.resource = resource; @@ -81,7 +105,9 @@ public LoadSingleTsFileNode( this.database = database; this.deleteAfterLoad = deleteAfterLoad; this.writePointCount = writePointCount; - this.needDecodeTsFile = needDecodeTsFile; + this.needDecodeTsFile = needDecode4TimeColumn; + this.tsFileContainsObjectColumn = tsFileContainsObjectColumn; + this.objectFileSearchRoot = objectFileSearchRoot; } public boolean isTsFileEmpty() { @@ -123,6 +149,18 @@ public boolean needDecodeTsFile( return needDecodeTsFile; } + public boolean isTsFileContainsObjectColumn() { + return tsFileContainsObjectColumn; + } + + public File getObjectFileSearchRoot() { + return objectFileSearchRoot; + } + + public void setObjectFileSearchRoot(final File objectFileSearchRoot) { + this.objectFileSearchRoot = objectFileSearchRoot; + } + private boolean isDispatchedToLocal(Set replicaSets) { if (replicaSets.size() > 1) { return false; @@ -262,6 +300,9 @@ public boolean equals(Object o) { && Objects.equals(isTableModel, loadSingleTsFileNode.isTableModel) && Objects.equals(database, loadSingleTsFileNode.database) && Objects.equals(needDecodeTsFile, loadSingleTsFileNode.needDecodeTsFile) + && Objects.equals( + tsFileContainsObjectColumn, loadSingleTsFileNode.tsFileContainsObjectColumn) + && Objects.equals(objectFileSearchRoot, loadSingleTsFileNode.objectFileSearchRoot) && Objects.equals(deleteAfterLoad, loadSingleTsFileNode.deleteAfterLoad) && Objects.equals(localRegionReplicaSet, loadSingleTsFileNode.localRegionReplicaSet); } @@ -274,6 +315,8 @@ public int hashCode() { isTableModel, database, needDecodeTsFile, + tsFileContainsObjectColumn, + objectFileSearchRoot, deleteAfterLoad, localRegionReplicaSet); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/ObjectNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/ObjectNode.java index 552d26e56319a..54d40015f9ce1 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/ObjectNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/ObjectNode.java @@ -70,6 +70,15 @@ public class ObjectNode extends SearchNode implements WALEntryValue { private boolean isGeneratedByRemoteConsensusLeader; + protected ObjectNode(final PlanNodeId planNodeId) { + super(planNodeId); + this.isEOF = false; + this.offset = 0; + this.filePath = null; + this.content = null; + this.contentLength = 0; + } + public ObjectNode(boolean isEOF, long offset, byte[] content, IObjectPath filePath) { super(new PlanNodeId("")); this.isEOF = isEOF; @@ -103,6 +112,10 @@ public void setFilePath(IObjectPath filePath) { this.filePath = filePath; } + public IObjectPath getFilePath() { + return filePath; + } + public String getFilePathString() { return filePath.toString(); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java index d43b828cced8f..4aaac36ff8868 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java @@ -93,6 +93,19 @@ public TsFileSplitter(File tsFile, TsFileDataConsumer consumer) { this.consumer = consumer; } + public TsFileSplitter( + File tsFile, TsFileDataConsumer consumer, boolean fileContainsObjectColumns) { + this(tsFile, consumer); + } + + public TsFileSplitter( + File tsFile, + TsFileDataConsumer consumer, + boolean fileContainsObjectColumns, + File objectFileSearchRoot) { + this(tsFile, consumer); + } + @SuppressWarnings({"squid:S3776", "squid:S6541"}) public void splitTsFileByDataPartition() throws IOException, LoadFileException, IllegalStateException {