This is an automated email from the ASF dual-hosted git repository. jackietien pushed a commit to branch NewTsFile in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 50b9da45446e553f89adcf1b106f8e391ce34474 Author: JackieTien97 <[email protected]> AuthorDate: Mon Nov 30 20:59:58 2020 +0800 have a good day --- .../apache/iotdb/tsfile/TsFileSequenceRead.java | 40 ++++--- .../iotdb/db/qp/physical/crud/InsertRowPlan.java | 4 +- .../db/query/reader/series/SeriesReaderTest.java | 19 ++- .../db/writelog/recover/SeqTsFileRecoverTest.java | 2 +- .../iotdb/tsfile/file/header/ChunkHeader.java | 16 ++- .../iotdb/tsfile/file/header/PageHeader.java | 19 ++- .../iotdb/tsfile/read/TsFileSequenceReader.java | 130 ++++++++++++++------- .../tsfile/read/reader/chunk/ChunkReader.java | 2 +- .../apache/iotdb/tsfile/write/TsFileWriter.java | 2 +- .../iotdb/tsfile/write/chunk/ChunkWriterImpl.java | 4 +- .../iotdb/tsfile/file/header/PageHeaderTest.java | 2 +- .../tsfile/read/TsFileSequenceReaderTest.java | 3 +- 12 files changed, 158 insertions(+), 85 deletions(-) diff --git a/example/tsfile/src/main/java/org/apache/iotdb/tsfile/TsFileSequenceRead.java b/example/tsfile/src/main/java/org/apache/iotdb/tsfile/TsFileSequenceRead.java index e9314fa..93ed9f0 100644 --- a/example/tsfile/src/main/java/org/apache/iotdb/tsfile/TsFileSequenceRead.java +++ b/example/tsfile/src/main/java/org/apache/iotdb/tsfile/TsFileSequenceRead.java @@ -41,18 +41,19 @@ public class TsFileSequenceRead { @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity warning public static void main(String[] args) throws IOException { - String filename = "test.tsfile"; + String filename = "/Users/jackietien/Desktop/1-1-1-after.tsfile"; if (args.length >= 1) { filename = args[0]; } try (TsFileSequenceReader reader = new TsFileSequenceReader(filename)) { - System.out.println("file length: " + FSFactoryProducer.getFSFactory().getFile(filename).length()); + System.out + .println("file length: " + FSFactoryProducer.getFSFactory().getFile(filename).length()); System.out.println("file magic head: " + reader.readHeadMagic()); System.out.println("file magic tail: " + reader.readTailMagic()); System.out.println("Level 1 metadata position: " + reader.getFileMetadataPos()); System.out.println("Level 1 metadata size: " + reader.getFileMetadataSize()); // Sequential reading of one ChunkGroup now follows this order: - // first SeriesChunks (headers and data) in one ChunkGroup, then the CHUNK_GROUP_FOOTER + // first the CHUNK_GROUP_HEADER, then SeriesChunks (headers and data) in one ChunkGroup // Because we do not know how many chunks a ChunkGroup may have, we should read one byte (the marker) ahead and // judge accordingly. reader.position((long) TSFileConfig.MAGIC_STRING.getBytes().length + 1); @@ -68,32 +69,39 @@ public class TsFileSequenceRead { ChunkHeader header = reader.readChunkHeader(marker); System.out.println("\tMeasurement: " + header.getMeasurementID()); Decoder defaultTimeDecoder = Decoder.getDecoderByType( - TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()), - TSDataType.INT64); + TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()), + TSDataType.INT64); Decoder valueDecoder = Decoder - .getDecoderByType(header.getEncodingType(), header.getDataType()); - for (int j = 0; j < header.getNumOfPages(); j++) { + .getDecoderByType(header.getEncodingType(), header.getDataType()); + int dataSize = header.getDataSize(); + while (dataSize > 0) { valueDecoder.reset(); System.out.println("\t\t[Page]\n \t\tPage head position: " + reader.position()); - PageHeader pageHeader = reader.readPageHeader(header.getDataType()); + PageHeader pageHeader = reader.readPageHeader(header.getDataType(), + header.getChunkType() == MetaMarker.CHUNK_HEADER); System.out.println("\t\tPage data position: " + reader.position()); - System.out.println("\t\tpoints in the page: " + pageHeader.getNumOfValues()); ByteBuffer pageData = reader.readPage(pageHeader, header.getCompressionType()); System.out - .println("\t\tUncompressed page data size: " + pageHeader.getUncompressedSize()); + .println("\t\tUncompressed page data size: " + pageHeader.getUncompressedSize()); PageReader reader1 = new PageReader(pageData, header.getDataType(), valueDecoder, - defaultTimeDecoder, null); + defaultTimeDecoder, null); BatchData batchData = reader1.getAllSatisfiedPageData(); + if (header.getChunkType() == MetaMarker.CHUNK_HEADER) { + System.out.println("\t\tpoints in the page: " + pageHeader.getNumOfValues()); + } else { + System.out.println("\t\tpoints in the page: " + batchData.length()); + } while (batchData.hasCurrent()) { System.out.println( - "\t\t\ttime, value: " + batchData.currentTime() + ", " + batchData - .currentValue()); + "\t\t\ttime, value: " + batchData.currentTime() + ", " + batchData + .currentValue()); batchData.next(); } + dataSize -= pageHeader.getSerializedPageSize(); } break; case MetaMarker.CHUNK_GROUP_HEADER: - System.out.println("Chunk Group Footer position: " + reader.position()); + System.out.println("Chunk Group Header position: " + reader.position()); ChunkGroupHeader chunkGroupHeader = reader.readChunkGroupHeader(); System.out.println("device: " + chunkGroupHeader.getDeviceID()); break; @@ -108,8 +116,8 @@ public class TsFileSequenceRead { System.out.println("[Metadata]"); for (String device : reader.getAllDevices()) { Map<String, List<ChunkMetadata>> seriesMetaData = reader.readChunkMetadataInDevice(device); - System.out.println(String - .format("\t[Device]Device %s, Number of Measurements %d", device, seriesMetaData.size())); + System.out.printf("\t[Device]Device %s, Number of Measurements %d%n", device, + seriesMetaData.size()); for (Map.Entry<String, List<ChunkMetadata>> serie : seriesMetaData.entrySet()) { System.out.println("\t\tMeasurement:" + serie.getKey()); for (ChunkMetadata chunkMetadata : serie.getValue()) { diff --git a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java index 47f5afa..8c624cd 100644 --- a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java +++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java @@ -50,7 +50,7 @@ import org.slf4j.LoggerFactory; public class InsertRowPlan extends InsertPlan { private static final Logger logger = LoggerFactory.getLogger(InsertRowPlan.class); - private static final short TYPE_RAW_STRING = -1; + private static final byte TYPE_RAW_STRING = -1; private long time; private Object[] values; @@ -357,7 +357,7 @@ public class InsertRowPlan extends InsertPlan { for (int i = 0; i < measurements.length; i++) { // types are not determined, the situation mainly occurs when the plan uses string values // and is forwarded to other nodes - short typeNum = ReadWriteIOUtils.readShort(buffer); + byte typeNum = (byte) ReadWriteIOUtils.read(buffer); if (typeNum == TYPE_RAW_STRING) { values[i] = ReadWriteIOUtils.readString(buffer); continue; diff --git a/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTest.java b/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTest.java index 9474a9b..4a9267f 100644 --- a/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTest.java +++ b/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTest.java @@ -19,13 +19,20 @@ package org.apache.iotdb.db.query.reader.series; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.fail; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; import org.apache.iotdb.db.exception.StorageEngineException; import org.apache.iotdb.db.exception.metadata.IllegalPathException; import org.apache.iotdb.db.exception.metadata.MetadataException; import org.apache.iotdb.db.metadata.PartialPath; import org.apache.iotdb.db.query.context.QueryContext; -import org.apache.iotdb.db.utils.TestOnly; import org.apache.iotdb.tsfile.exception.write.WriteProcessException; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import org.apache.iotdb.tsfile.read.TimeValuePair; @@ -37,15 +44,6 @@ import org.junit.After; import org.junit.Before; import org.junit.Test; -import java.io.IOException; -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.fail; - public class SeriesReaderTest { private static final String SERIES_READER_TEST_SG = "root.seriesReaderTest"; @@ -144,7 +142,6 @@ public class SeriesReaderTest { long expectedTime = 499; while (pointReader.hasNextTimeValuePair()) { TimeValuePair timeValuePair = pointReader.nextTimeValuePair(); - System.out.println(timeValuePair); assertEquals(expectedTime, timeValuePair.getTimestamp()); int value = timeValuePair.getValue().getInt(); if (expectedTime < 200) { diff --git a/server/src/test/java/org/apache/iotdb/db/writelog/recover/SeqTsFileRecoverTest.java b/server/src/test/java/org/apache/iotdb/db/writelog/recover/SeqTsFileRecoverTest.java index d550f60..5a98ab2 100644 --- a/server/src/test/java/org/apache/iotdb/db/writelog/recover/SeqTsFileRecoverTest.java +++ b/server/src/test/java/org/apache/iotdb/db/writelog/recover/SeqTsFileRecoverTest.java @@ -73,7 +73,6 @@ public class SeqTsFileRecoverTest { private WriteLogNode node; private String logNodePrefix = TestConstant.BASE_OUTPUT_PATH.concat("testRecover"); - private String storageGroup = "target"; private TsFileResource resource; private VersionController versionController = new VersionController() { private int i; @@ -136,6 +135,7 @@ public class SeqTsFileRecoverTest { } } writer.flushAllChunkGroups(); + writer.writeVersion(0); writer.getIOWriter().close(); node = MultiFileLogNodeManager.getInstance().getNode(logNodePrefix + tsF.getName()); diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/ChunkHeader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/ChunkHeader.java index 62267e6..04942d8 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/ChunkHeader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/ChunkHeader.java @@ -81,9 +81,10 @@ public class ChunkHeader { */ public static int getSerializedSize(String measurementID, int dataSize) { int measurementIdLength = measurementID.getBytes(TSFileConfig.STRING_CHARSET).length; - return ReadWriteForEncodingUtils.varIntSize(measurementIdLength) // measurementID length + return Byte.BYTES // chunkType + + ReadWriteForEncodingUtils.varIntSize(measurementIdLength) // measurementID length + measurementIdLength // measurementID - + ReadWriteForEncodingUtils.varIntSize(dataSize) // dataSize + + ReadWriteForEncodingUtils.uVarIntSize(dataSize) // dataSize + TSDataType.getSerializedSize() // dataType + CompressionType.getSerializedSize() // compressionType + TSEncoding.getSerializedSize(); // encodingType @@ -96,9 +97,10 @@ public class ChunkHeader { public static int getSerializedSize(String measurementID) { int measurementIdLength = measurementID.getBytes(TSFileConfig.STRING_CHARSET).length; - return ReadWriteForEncodingUtils.varIntSize(measurementIdLength) // measurementID length + return Byte.BYTES // chunkType + + ReadWriteForEncodingUtils.varIntSize(measurementIdLength) // measurementID length + measurementIdLength // measurementID - + Integer.BYTES + 1 // varInr dataSize + + Integer.BYTES + 1 // uVarInt dataSize + TSDataType.getSerializedSize() // dataType + CompressionType.getSerializedSize() // compressionType + TSEncoding.getSerializedSize(); // encodingType @@ -142,7 +144,7 @@ public class ChunkHeader { CompressionType type = ReadWriteIOUtils.readCompressionType(buffer); TSEncoding encoding = ReadWriteIOUtils.readEncoding(buffer); chunkHeaderSize = - chunkHeaderSize - Integer.BYTES + ReadWriteForEncodingUtils.varIntSize(dataSize); + chunkHeaderSize - Integer.BYTES - 1 + ReadWriteForEncodingUtils.uVarIntSize(dataSize); return new ChunkHeader(chunkType, measurementID, dataSize, chunkHeaderSize, dataType, type, encoding); } @@ -227,4 +229,8 @@ public class ChunkHeader { public byte getChunkType() { return chunkType; } + + public void increasePageNums(int i) { + numOfPages += i; + } } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/PageHeader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/PageHeader.java index 2c0acf9..a990430 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/PageHeader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/header/PageHeader.java @@ -47,11 +47,14 @@ public class PageHeader { return 2 * (Integer.BYTES + 1); // uncompressedSize, compressedSize } - public static PageHeader deserializeFrom(InputStream inputStream, TSDataType dataType) - throws IOException { + public static PageHeader deserializeFrom(InputStream inputStream, TSDataType dataType, + boolean hasStatistic) throws IOException { int uncompressedSize = ReadWriteForEncodingUtils.readUnsignedVarInt(inputStream); int compressedSize = ReadWriteForEncodingUtils.readUnsignedVarInt(inputStream); - Statistics statistics = Statistics.deserialize(inputStream, dataType); + Statistics statistics = null; + if (hasStatistic) { + statistics = Statistics.deserialize(inputStream, dataType); + } return new PageHeader(uncompressedSize, compressedSize, statistics); } @@ -119,4 +122,14 @@ public class PageHeader { public void setModified(boolean modified) { this.modified = modified; } + + /** + * max page header size without statistics + */ + public int getSerializedPageSize() { + return ReadWriteForEncodingUtils.uVarIntSize(uncompressedSize) + + ReadWriteForEncodingUtils.uVarIntSize(compressedSize) + + (statistics == null ? 0 : statistics.getSerializedSize()) // page header + + compressedSize; // page data + } } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java index 5c7338d..c3d2c5a 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java @@ -38,6 +38,7 @@ import java.util.stream.Collectors; import org.apache.iotdb.tsfile.common.conf.TSFileConfig; import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor; import org.apache.iotdb.tsfile.compress.IUnCompressor; +import org.apache.iotdb.tsfile.encoding.decoder.Decoder; import org.apache.iotdb.tsfile.file.MetaMarker; import org.apache.iotdb.tsfile.file.header.ChunkGroupHeader; import org.apache.iotdb.tsfile.file.header.ChunkHeader; @@ -51,12 +52,15 @@ import org.apache.iotdb.tsfile.file.metadata.TsFileMetadata; import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType; import org.apache.iotdb.tsfile.file.metadata.enums.MetadataIndexNodeType; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; +import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding; import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics; import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer; +import org.apache.iotdb.tsfile.read.common.BatchData; import org.apache.iotdb.tsfile.read.common.Chunk; import org.apache.iotdb.tsfile.read.common.Path; import org.apache.iotdb.tsfile.read.controller.MetadataQuerierByFileImpl; import org.apache.iotdb.tsfile.read.reader.TsFileInput; +import org.apache.iotdb.tsfile.read.reader.page.PageReader; import org.apache.iotdb.tsfile.utils.BloomFilter; import org.apache.iotdb.tsfile.utils.Pair; import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; @@ -75,7 +79,6 @@ public class TsFileSequenceReader implements AutoCloseable { private long fileMetadataPos; private int fileMetadataSize; private ByteBuffer markerBuffer = ByteBuffer.allocate(Byte.BYTES); - private int totalChunkNum; private TsFileMetadata tsFileMetaData; // device -> measurement -> TimeseriesMetadata private Map<String, Map<String, TimeseriesMetadata>> cachedDeviceMetadata = new ConcurrentHashMap<>(); @@ -206,10 +209,8 @@ public class TsFileSequenceReader implements AutoCloseable { * whether the file is a complete TsFile: only if the head magic and tail magic string exists. */ public boolean isComplete() throws IOException { - return tsFileInput.size() >= TSFileConfig.MAGIC_STRING.getBytes().length * 2 - + TSFileConfig.VERSION_NUMBER_V2.getBytes().length - && (readTailMagic().equals(readHeadMagic()) || readTailMagic() - .equals(TSFileConfig.VERSION_NUMBER_V1)); + return tsFileInput.size() >= TSFileConfig.MAGIC_STRING.getBytes().length * 2 + Byte.BYTES + && (readTailMagic().equals(readHeadMagic())); } /** @@ -763,8 +764,8 @@ public class TsFileSequenceReader implements AutoCloseable { * * @param type given tsfile data type */ - public PageHeader readPageHeader(TSDataType type) throws IOException { - return PageHeader.deserializeFrom(tsFileInput.wrapAsInputStream(), type); + public PageHeader readPageHeader(TSDataType type, boolean hasStatistic) throws IOException { + return PageHeader.deserializeFrom(tsFileInput.wrapAsInputStream(), type, hasStatistic); } public long position() throws IOException { @@ -902,11 +903,9 @@ public class TsFileSequenceReader implements AutoCloseable { long fileOffsetOfChunk; // ChunkMetadata of current ChunkGroup - List<ChunkMetadata> chunkMetadataList = null; - String deviceID; + List<ChunkMetadata> chunkMetadataList = new ArrayList<>(); - int headerLength = TSFileConfig.MAGIC_STRING.getBytes().length + TSFileConfig.VERSION_NUMBER_V2 - .getBytes().length; + int headerLength = TSFileConfig.MAGIC_STRING.getBytes().length + Byte.BYTES; if (fileSize < headerLength) { return TsFileCheckStatus.INCOMPATIBLE_FILE; } @@ -924,22 +923,16 @@ public class TsFileSequenceReader implements AutoCloseable { return TsFileCheckStatus.COMPLETE_FILE; } } - boolean newChunkGroup = true; // not a complete file, we will recover it... long truncatedSize = headerLength; byte marker; - int chunkCnt = 0; + String lastDeviceId = null; List<MeasurementSchema> measurementSchemaList = new ArrayList<>(); try { while ((marker = this.readMarker()) != MetaMarker.SEPARATOR) { switch (marker) { case MetaMarker.CHUNK_HEADER: case MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER: - // this is the first chunk of a new ChunkGroup. - if (newChunkGroup) { - newChunkGroup = false; - chunkMetadataList = new ArrayList<>(); - } fileOffsetOfChunk = this.position() - 1; // if there is something wrong with a chunk, we will drop the whole ChunkGroup // as different chunks may be created by the same insertions(sqls), and partial @@ -952,37 +945,96 @@ public class TsFileSequenceReader implements AutoCloseable { measurementSchemaList.add(measurementSchema); dataType = chunkHeader.getDataType(); Statistics<?> chunkStatistics = Statistics.getStatsByType(dataType); - for (int j = 0; j < chunkHeader.getNumOfPages(); j++) { - // a new Page - PageHeader pageHeader = this.readPageHeader(chunkHeader.getDataType()); - chunkStatistics.mergeStatistics(pageHeader.getStatistics()); - this.skipPageData(pageHeader); + int dataSize = chunkHeader.getDataSize(); + if (chunkHeader.getChunkType() == MetaMarker.CHUNK_HEADER) { + while (dataSize > 0) { + // a new Page + PageHeader pageHeader = this.readPageHeader(chunkHeader.getDataType(), true); + chunkStatistics.mergeStatistics(pageHeader.getStatistics()); + this.skipPageData(pageHeader); + dataSize -= pageHeader.getSerializedPageSize(); + chunkHeader.increasePageNums(1); + } + } else { + // only one page without statistic, we need to iterate each point to generate statistic + PageHeader pageHeader = this.readPageHeader(chunkHeader.getDataType(), false); + Decoder valueDecoder = Decoder + .getDecoderByType(chunkHeader.getEncodingType(), chunkHeader.getDataType()); + ByteBuffer pageData = readPage(pageHeader, chunkHeader.getCompressionType()); + Decoder timeDecoder = Decoder.getDecoderByType( + TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()), + TSDataType.INT64); + PageReader reader = new PageReader(pageHeader, pageData, chunkHeader.getDataType(), + valueDecoder, timeDecoder, null); + BatchData batchData = reader.getAllSatisfiedPageData(); + while (batchData.hasCurrent()) { + switch (dataType) { + case INT32: + chunkStatistics.update(batchData.currentTime(), batchData.getInt()); + break; + case INT64: + chunkStatistics.update(batchData.currentTime(), batchData.getLong()); + break; + case FLOAT: + chunkStatistics.update(batchData.currentTime(), batchData.getFloat()); + break; + case DOUBLE: + chunkStatistics.update(batchData.currentTime(), batchData.getDouble()); + break; + case BOOLEAN: + chunkStatistics.update(batchData.currentTime(), batchData.getBoolean()); + break; + case TEXT: + chunkStatistics.update(batchData.currentTime(), batchData.getBinary()); + break; + default: + throw new IOException("Unexpected type " + dataType); + } + batchData.next(); + } + chunkHeader.increasePageNums(1); } currentChunk = new ChunkMetadata(measurementID, dataType, fileOffsetOfChunk, chunkStatistics); chunkMetadataList.add(currentChunk); - chunkCnt++; break; case MetaMarker.CHUNK_GROUP_HEADER: - // this is a chunk group + if (lastDeviceId != null) { + // schema of last chunk group + if (newSchema != null) { + for (MeasurementSchema tsSchema : measurementSchemaList) { + newSchema + .putIfAbsent(new Path(lastDeviceId, tsSchema.getMeasurementId()), tsSchema); + } + } + measurementSchemaList = new ArrayList<>(); + // last chunk group Metadata + chunkGroupMetadataList.add(new ChunkGroupMetadata(lastDeviceId, chunkMetadataList)); + } // if there is something wrong with the ChunkGroup Footer, we will drop this ChunkGroup // because we can not guarantee the correctness of the deviceId. + truncatedSize = this.position() - 1; + // this is a chunk group + chunkMetadataList = new ArrayList<>(); ChunkGroupHeader chunkGroupHeader = this.readChunkGroupHeader(); - deviceID = chunkGroupHeader.getDeviceID(); - if (newSchema != null) { - for (MeasurementSchema tsSchema : measurementSchemaList) { - newSchema.putIfAbsent(new Path(deviceID, tsSchema.getMeasurementId()), tsSchema); + lastDeviceId = chunkGroupHeader.getDeviceID(); + break; + case MetaMarker.VERSION: + if (lastDeviceId != null) { + // schema of last chunk group + if (newSchema != null) { + for (MeasurementSchema tsSchema : measurementSchemaList) { + newSchema + .putIfAbsent(new Path(lastDeviceId, tsSchema.getMeasurementId()), tsSchema); + } } + measurementSchemaList = new ArrayList<>(); + // last chunk group Metadata + chunkGroupMetadataList.add(new ChunkGroupMetadata(lastDeviceId, chunkMetadataList)); + lastDeviceId = null; } - chunkGroupMetadataList.add(new ChunkGroupMetadata(deviceID, chunkMetadataList)); - newChunkGroup = true; - truncatedSize = this.position(); - totalChunkNum += chunkCnt; - chunkCnt = 0; - measurementSchemaList = new ArrayList<>(); - break; - case MetaMarker.VERSION: + chunkMetadataList = new ArrayList<>(); long version = readVersion(); versionInfo.add(new Pair<>(position(), version)); truncatedSize = this.position(); @@ -1004,10 +1056,6 @@ public class TsFileSequenceReader implements AutoCloseable { return truncatedSize; } - public int getTotalChunkNum() { - return totalChunkNum; - } - /** * get ChunkMetaDatas of given path * diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/ChunkReader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/ChunkReader.java index 729b8f7..2250307 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/ChunkReader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/ChunkReader.java @@ -45,7 +45,7 @@ public class ChunkReader implements IChunkReader { private ChunkHeader chunkHeader; private ByteBuffer chunkDataBuffer; private IUnCompressor unCompressor; - private Decoder timeDecoder = Decoder.getDecoderByType( + private final Decoder timeDecoder = Decoder.getDecoderByType( TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()), TSDataType.INT64); diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/TsFileWriter.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/TsFileWriter.java index 5346847..6b890dd 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/TsFileWriter.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/TsFileWriter.java @@ -314,10 +314,10 @@ public class TsFileWriter implements AutoCloseable { public boolean flushAllChunkGroups() throws IOException { if (recordCount > 0) { for (Map.Entry<String, IChunkGroupWriter> entry : groupWriters.entrySet()) { - long pos = fileWriter.getPos(); String deviceId = entry.getKey(); IChunkGroupWriter groupWriter = entry.getValue(); fileWriter.startChunkGroup(deviceId); + long pos = fileWriter.getPos(); long dataSize = groupWriter.flushToFileWriter(fileWriter); if (fileWriter.getPos() - pos != dataSize) { throw new IOException( diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java index 107f026..44c2dae 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java @@ -213,9 +213,9 @@ public class ChunkWriterImpl implements IChunkWriter { } else if (numOfPages == 1) { // put the firstPageStatistics into pageBuffer byte[] b = pageBuffer.toByteArray(); pageBuffer.reset(); - pageBuffer.write(b, 0, sizeWithoutStatistic); + pageBuffer.write(b, 0, this.sizeWithoutStatistic); firstPageStatistics.serialize(pageBuffer); - pageBuffer.write(b, sizeWithoutStatistic, b.length - sizeWithoutStatistic); + pageBuffer.write(b, this.sizeWithoutStatistic, b.length - this.sizeWithoutStatistic); firstPageStatistics = null; } diff --git a/tsfile/src/test/java/org/apache/iotdb/tsfile/file/header/PageHeaderTest.java b/tsfile/src/test/java/org/apache/iotdb/tsfile/file/header/PageHeaderTest.java index 3114159..ef87b31 100644 --- a/tsfile/src/test/java/org/apache/iotdb/tsfile/file/header/PageHeaderTest.java +++ b/tsfile/src/test/java/org/apache/iotdb/tsfile/file/header/PageHeaderTest.java @@ -70,7 +70,7 @@ public class PageHeaderTest { PageHeader header = null; try { fis = new FileInputStream(new File(PATH)); - header = PageHeader.deserializeFrom(fis, DATA_TYPE); + header = PageHeader.deserializeFrom(fis, DATA_TYPE, true); return header; } catch (IOException e) { e.printStackTrace(); diff --git a/tsfile/src/test/java/org/apache/iotdb/tsfile/read/TsFileSequenceReaderTest.java b/tsfile/src/test/java/org/apache/iotdb/tsfile/read/TsFileSequenceReaderTest.java index c9931fc..4e8b7e0 100644 --- a/tsfile/src/test/java/org/apache/iotdb/tsfile/read/TsFileSequenceReaderTest.java +++ b/tsfile/src/test/java/org/apache/iotdb/tsfile/read/TsFileSequenceReaderTest.java @@ -75,7 +75,8 @@ public class TsFileSequenceReaderTest { case MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER: ChunkHeader header = reader.readChunkHeader(marker); for (int j = 0; j < header.getNumOfPages(); j++) { - PageHeader pageHeader = reader.readPageHeader(header.getDataType()); + PageHeader pageHeader = reader.readPageHeader(header.getDataType(), + header.getChunkType() == MetaMarker.CHUNK_HEADER); reader.readPage(pageHeader, header.getCompressionType()); } break;
