From d754395e3f2b9296a5fd51b8112daaf9211ea774 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Mon, 3 Aug 2026 16:28:46 +0800 Subject: [PATCH 1/3] fix(load): use buffered input when splitting TsFiles --- .../iotdb/db/storageengine/load/splitter/TsFileSplitter.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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..ddd0878658205 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()))) { getAllModification(deletions); if (!checkMagic(reader)) { From ea4c171f46e548ecab002fdf74fb810875f554c8 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Tue, 4 Aug 2026 11:30:27 +0800 Subject: [PATCH 2/3] fixed --- .../load/splitter/TsFileSplitter.java | 2 +- .../load/TsFileSplitterTest.java | 83 +++++++++++++++++++ 2 files changed, 84 insertions(+), 1 deletion(-) 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 ddd0878658205..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 @@ -98,7 +98,7 @@ public TsFileSplitter(File tsFile, TsFileDataConsumer consumer) { public void splitTsFileByDataPartition() throws IOException, LoadFileException, IllegalStateException { try (TsFileSequenceReader reader = - new TsFileSequenceReader(new BufferedTsFileInput(tsFile.toPath()))) { + 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..d9e5b2ba1db94 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,13 +19,24 @@ 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.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.file.metadata.statistics.Statistics; import org.apache.tsfile.read.TsFileSequenceReader; +import org.apache.tsfile.utils.PublicBAOS; +import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.write.chunk.AlignedChunkWriterImpl; import org.apache.tsfile.write.schema.IMeasurementSchema; import org.apache.tsfile.write.schema.MeasurementSchema; @@ -38,6 +49,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 +59,76 @@ 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(); + + try { + try (final TsFileSequenceReader reader = + new TsFileSequenceReader(sourceTsFile.getAbsolutePath())) { + Assert.assertEquals(2, reader.getAllTimeseriesMetadata(true).size()); + } + + // Verify the buffered reader initializes the v3 deserialize configuration before reading + // metadata. + new TsFileSplitter(sourceTsFile, tsFileData -> true).splitTsFileByDataPartition(); + } finally { + Assert.assertTrue(sourceTsFile.delete()); + } + } + + private File constructV3TsFile() throws IOException { + final File tsFile = Files.createTempFile("v3-tsfile-splitter", ".tsfile").toFile(); + final ByteArrayOutputStream outputStream = new ByteArrayOutputStream(); + outputStream.write(TSFileConfig.MAGIC_STRING.getBytes()); + outputStream.write(TSFileConfig.VERSION_NUMBER_V3); + + final long metaOffset = outputStream.size(); + outputStream.write(MetaMarker.SEPARATOR); + + final MetadataIndexNode deviceIndexNode = + new MetadataIndexNode(MetadataIndexNodeType.LEAF_DEVICE); + for (final String device : Arrays.asList("root.sg.d1", "root.sg.d2")) { + final long timeseriesMetadataOffset = outputStream.size(); + writeV3TimeseriesMetadata(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); + + deviceIndexNode.addEntry( + new DeviceMetadataIndexEntry(new PlainDeviceID(device), 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; + } + + private void writeV3TimeseriesMetadata(final ByteArrayOutputStream outputStream) + throws IOException { + final Statistics statistics = Statistics.getStatsByType(TSDataType.INT32); + statistics.update(1, 1); + + final TimeseriesMetadata timeseriesMetadata = new TimeseriesMetadata(); + timeseriesMetadata.setTimeSeriesMetadataType((byte) 0); + timeseriesMetadata.setMeasurementId("s1"); + timeseriesMetadata.setTsDataType(TSDataType.INT32); + timeseriesMetadata.setDataSizeOfChunkMetaDataList(0); + timeseriesMetadata.setStatistics(statistics); + timeseriesMetadata.setChunkMetadataListBuffer(new PublicBAOS()); + timeseriesMetadata.serializeTo(outputStream); + } + @Test public void testSplitTableTimeOnlyAlignedChunk() throws Exception { final File sourceTsFile = new File("split-table-time-only-source.tsfile"); From 03397b41be6eeaad0d78f97f5d1465bb0645f7c8 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Tue, 4 Aug 2026 14:42:26 +0800 Subject: [PATCH 3/3] test(load): assert v3 splitter output --- .../load/TsFileSplitterTest.java | 80 ++++++++++--------- 1 file changed, 44 insertions(+), 36 deletions(-) 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 d9e5b2ba1db94..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 @@ -25,6 +25,7 @@ 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; @@ -33,15 +34,15 @@ 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.file.metadata.statistics.Statistics; import org.apache.tsfile.read.TsFileSequenceReader; -import org.apache.tsfile.utils.PublicBAOS; 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; @@ -63,16 +64,26 @@ public class TsFileSplitterTest { @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(2, reader.getAllTimeseriesMetadata(true).size()); + Assert.assertEquals(1, reader.getAllTimeseriesMetadata(true).size()); } // Verify the buffered reader initializes the v3 deserialize configuration before reading // metadata. - new TsFileSplitter(sourceTsFile, tsFileData -> true).splitTsFileByDataPartition(); + 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()); } @@ -80,30 +91,42 @@ public void testSplitV3TsFile() throws Exception { 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(TSFileConfig.MAGIC_STRING.getBytes()); - outputStream.write(TSFileConfig.VERSION_NUMBER_V3); + 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); - for (final String device : Arrays.asList("root.sg.d1", "root.sg.d2")) { - final long timeseriesMetadataOffset = outputStream.size(); - writeV3TimeseriesMetadata(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); - - deviceIndexNode.addEntry( - new DeviceMetadataIndexEntry(new PlainDeviceID(device), measurementIndexOffset)); - } + deviceIndexNode.addEntry(new DeviceMetadataIndexEntry(deviceID, measurementIndexOffset)); deviceIndexNode.setEndOffset(outputStream.size()); int metadataSize = deviceIndexNode.serializeTo(outputStream); @@ -114,21 +137,6 @@ private File constructV3TsFile() throws IOException { return tsFile; } - private void writeV3TimeseriesMetadata(final ByteArrayOutputStream outputStream) - throws IOException { - final Statistics statistics = Statistics.getStatsByType(TSDataType.INT32); - statistics.update(1, 1); - - final TimeseriesMetadata timeseriesMetadata = new TimeseriesMetadata(); - timeseriesMetadata.setTimeSeriesMetadataType((byte) 0); - timeseriesMetadata.setMeasurementId("s1"); - timeseriesMetadata.setTsDataType(TSDataType.INT32); - timeseriesMetadata.setDataSizeOfChunkMetaDataList(0); - timeseriesMetadata.setStatistics(statistics); - timeseriesMetadata.setChunkMetadataListBuffer(new PublicBAOS()); - timeseriesMetadata.serializeTo(outputStream); - } - @Test public void testSplitTableTimeOnlyAlignedChunk() throws Exception { final File sourceTsFile = new File("split-table-time-only-source.tsfile");