Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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)) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,25 +19,39 @@

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;

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;
Expand All @@ -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<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");
Expand Down
Loading