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 f8ba40cbe687b..60f382e0cd9bd 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 @@ -44,6 +44,7 @@ import org.apache.tsfile.file.metadata.enums.TSEncoding; import org.apache.tsfile.read.TsFileSequenceReader; import org.apache.tsfile.read.common.BatchData; +import org.apache.tsfile.read.reader.BufferedTsFileInput; import org.apache.tsfile.read.reader.page.PageReader; import org.apache.tsfile.read.reader.page.TimePageReader; import org.apache.tsfile.read.reader.page.ValuePageReader; @@ -96,7 +97,8 @@ public TsFileSplitter(File tsFile, TsFileDataConsumer consumer) { @SuppressWarnings({"squid:S3776", "squid:S6541"}) public void splitTsFileByDataPartition() throws IOException, LoadFileException, IllegalStateException { - try (TsFileSequenceReader reader = new TsFileSequenceReader(tsFile.getAbsolutePath())) { + try (TsFileSequenceReader reader = + new TsFileSequenceReader(new BufferedTsFileInput(tsFile.toPath()), true, false, null)) { getAllModification(deletions); if (!checkMagic(reader)) { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/TsFileSplitterTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/TsFileSplitterTest.java index 6610880567e90..dd923c032272d 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/TsFileSplitterTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/TsFileSplitterTest.java @@ -19,18 +19,30 @@ package org.apache.iotdb.db.storageengine.load.splitter; +import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.enums.ColumnCategory; import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.MetaMarker; import org.apache.tsfile.file.metadata.AbstractAlignedChunkMetadata; +import org.apache.tsfile.file.metadata.DeviceMetadataIndexEntry; +import org.apache.tsfile.file.metadata.IChunkMetadata; import org.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.file.metadata.MeasurementMetadataIndexEntry; +import org.apache.tsfile.file.metadata.MetadataIndexNode; +import org.apache.tsfile.file.metadata.PlainDeviceID; import org.apache.tsfile.file.metadata.StringArrayDeviceID; import org.apache.tsfile.file.metadata.TableSchema; +import org.apache.tsfile.file.metadata.TimeseriesMetadata; +import org.apache.tsfile.file.metadata.enums.MetadataIndexNodeType; import org.apache.tsfile.read.TsFileSequenceReader; +import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.write.chunk.AlignedChunkWriterImpl; +import org.apache.tsfile.write.chunk.ChunkWriterImpl; import org.apache.tsfile.write.schema.IMeasurementSchema; import org.apache.tsfile.write.schema.MeasurementSchema; import org.apache.tsfile.write.schema.Schema; import org.apache.tsfile.write.writer.TsFileIOWriter; +import org.apache.tsfile.write.writer.tsmiterator.TSMIterator; import org.junit.Assert; import org.junit.Test; @@ -38,6 +50,8 @@ import java.io.ByteArrayOutputStream; import java.io.DataOutputStream; import java.io.File; +import java.io.IOException; +import java.nio.file.Files; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -46,6 +60,83 @@ public class TsFileSplitterTest { + // Verify the splitter initializes the v3 deserialize configuration for a valid v3 TsFile. + @Test + public void testSplitV3TsFile() throws Exception { + final File sourceTsFile = constructV3TsFile(); + final List chunkDataList = new ArrayList<>(); + + try { + try (final TsFileSequenceReader reader = + new TsFileSequenceReader(sourceTsFile.getAbsolutePath())) { + Assert.assertEquals(1, reader.getAllTimeseriesMetadata(true).size()); + } + + // Verify the buffered reader initializes the v3 deserialize configuration before reading + // metadata. + new TsFileSplitter( + sourceTsFile, + tsFileData -> { + if (tsFileData instanceof ChunkData) { + chunkDataList.add((ChunkData) tsFileData); + } + return true; + }) + .splitTsFileByDataPartition(); + Assert.assertEquals(1, chunkDataList.size()); + } finally { + Assert.assertTrue(sourceTsFile.delete()); + } + } + + private File constructV3TsFile() throws IOException { + final File tsFile = Files.createTempFile("v3-tsfile-splitter", ".tsfile").toFile(); + final IDeviceID deviceID = new PlainDeviceID("root.sg.d1"); + final TimeseriesMetadata timeseriesMetadata; + try (final TsFileIOWriter writer = new TsFileIOWriter(tsFile)) { + writer.startChunkGroup(deviceID); + final ChunkWriterImpl chunkWriter = + new ChunkWriterImpl(new MeasurementSchema("s1", TSDataType.INT32)); + chunkWriter.write(1, 1); + chunkWriter.writeToFileWriter(writer); + writer.endChunkGroup(); + + final List chunkMetadataList = + new ArrayList(writer.getDeviceChunkMetadataMap().get(deviceID)); + timeseriesMetadata = TSMIterator.constructOneTimeseriesMetadata("s1", chunkMetadataList); + } + + final byte[] v3TsFileData = Files.readAllBytes(tsFile.toPath()); + v3TsFileData[TSFileConfig.MAGIC_STRING.getBytes().length] = TSFileConfig.VERSION_NUMBER_V3; + final ByteArrayOutputStream outputStream = new ByteArrayOutputStream(); + outputStream.write(v3TsFileData); + + final long metaOffset = outputStream.size(); + outputStream.write(MetaMarker.SEPARATOR); + final long timeseriesMetadataOffset = outputStream.size(); + timeseriesMetadata.serializeTo(outputStream); + + final long measurementIndexOffset = outputStream.size(); + final MetadataIndexNode measurementIndexNode = + new MetadataIndexNode(MetadataIndexNodeType.LEAF_MEASUREMENT); + measurementIndexNode.addEntry( + new MeasurementMetadataIndexEntry("s1", timeseriesMetadataOffset)); + measurementIndexNode.setEndOffset(measurementIndexOffset); + measurementIndexNode.serializeTo(outputStream); + + final MetadataIndexNode deviceIndexNode = + new MetadataIndexNode(MetadataIndexNodeType.LEAF_DEVICE); + deviceIndexNode.addEntry(new DeviceMetadataIndexEntry(deviceID, measurementIndexOffset)); + deviceIndexNode.setEndOffset(outputStream.size()); + + int metadataSize = deviceIndexNode.serializeTo(outputStream); + metadataSize += ReadWriteIOUtils.write(metaOffset, outputStream); + ReadWriteIOUtils.write(metadataSize, outputStream); + outputStream.write(TSFileConfig.MAGIC_STRING.getBytes()); + Files.write(tsFile.toPath(), outputStream.toByteArray()); + return tsFile; + } + @Test public void testSplitTableTimeOnlyAlignedChunk() throws Exception { final File sourceTsFile = new File("split-table-time-only-source.tsfile");