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 7e2c453fbd1 fix(load): use buffered input when splitting TsFiles
(#18388)
7e2c453fbd1 is described below
commit 7e2c453fbd1880e977c4da50395881a36edc1898
Author: Zhenyu Luo <[email protected]>
AuthorDate: Fri Aug 7 09:28:22 2026 +0800
fix(load): use buffered input when splitting TsFiles (#18388)
* fix(load): use buffered input when splitting TsFiles
* fixed
* test(load): assert v3 splitter output
---
.../load/splitter/TsFileSplitter.java | 4 +-
.../db/storageengine/load/TsFileSplitterTest.java | 91 ++++++++++++++++++++++
2 files changed, 94 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 f8ba40cbe68..60f382e0cd9 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.TimeseriesMetadata;
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 class TsFileSplitter {
@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 6610880567e..dd923c03227 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.ByteArrayInputStream;
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 @@ import java.util.Objects;
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<ChunkData> 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<IChunkMetadata> chunkMetadataList =
+ new
ArrayList<IChunkMetadata>(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");