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 57bc1b69098b5759ec4d27c996db5bca4b8585d8 Author: JackieTien97 <[email protected]> AuthorDate: Fri Nov 27 20:32:05 2020 +0800 some changes --- .../apache/iotdb/db/engine/cache/ChunkCache.java | 5 +-- .../iotdb/db/query/control/FileReaderManager.java | 2 +- .../org/apache/iotdb/db/utils/FileLoaderUtils.java | 9 +++--- .../org/apache/iotdb/db/utils/UpgradeUtils.java | 2 +- .../iotdb/tsfile/file/header/ChunkHeader.java | 4 +++ .../iotdb/tsfile/file/header/PageHeader.java | 6 ++++ .../iotdb/tsfile/read/TsFileSequenceReader.java | 17 +++++----- .../org/apache/iotdb/tsfile/read/common/Chunk.java | 9 +++++- .../read/controller/CachedChunkLoaderImpl.java | 3 +- .../tsfile/read/reader/chunk/ChunkReader.java | 37 ++++++++++++---------- .../iotdb/tsfile/utils/ReadWriteIOUtils.java | 2 +- .../iotdb/tsfile/write/writer/TsFileIOWriter.java | 3 +- 12 files changed, 61 insertions(+), 38 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/engine/cache/ChunkCache.java b/server/src/main/java/org/apache/iotdb/db/engine/cache/ChunkCache.java index 6fccb8b..c621a26 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/cache/ChunkCache.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/cache/ChunkCache.java @@ -87,7 +87,7 @@ public class ChunkCache { if (!CACHE_ENABLE) { Chunk chunk = reader.readMemChunk(chunkMetaData); return new Chunk(chunk.getHeader(), chunk.getData().duplicate(), - chunk.getDeleteIntervalList()); + chunk.getDeleteIntervalList(), chunkMetaData.getStatistics()); } cacheRequestNum.incrementAndGet(); @@ -121,7 +121,8 @@ public class ChunkCache { if (config.isDebugOn()) { DEBUG_LOGGER.info("get chunk from cache whose meta data is: " + chunkMetaData); } - return new Chunk(chunk.getHeader(), chunk.getData().duplicate(), chunk.getDeleteIntervalList()); + return new Chunk(chunk.getHeader(), chunk.getData().duplicate(), chunk.getDeleteIntervalList(), + chunkMetaData.getStatistics()); } private void printCacheLog(boolean isHit) { diff --git a/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java b/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java index 9a9eb18..4e99292 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java +++ b/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java @@ -167,7 +167,7 @@ public class FileReaderManager implements IService { else { tsFileReader = new TsFileSequenceReader(filePath); switch (tsFileReader.readVersionNumber()) { - case TSFileConfig.VERSION_NUMBER_V2: + case TSFileConfig.VERSION_NUMBER: break; default: throw new IOException("The version of this TsFile is not corrent. "); diff --git a/server/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java b/server/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java index b3af734..a42d80f 100644 --- a/server/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java +++ b/server/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java @@ -80,12 +80,13 @@ public class FileLoaderUtils { } /** - * @param resource TsFile + * @param resource TsFile * @param seriesPath Timeseries path * @param allSensors measurements queried at the same time of this device - * @param filter any filter, only used to check time range + * @param filter any filter, only used to check time range */ - public static TimeseriesMetadata loadTimeSeriesMetadata(TsFileResource resource, PartialPath seriesPath, + public static TimeseriesMetadata loadTimeSeriesMetadata(TsFileResource resource, + PartialPath seriesPath, QueryContext context, Filter filter, Set<String> allSensors) throws IOException { TimeseriesMetadata timeSeriesMetadata; if (resource.isClosed()) { @@ -138,7 +139,7 @@ public class FileLoaderUtils { * load all page readers in one chunk that satisfying the timeFilter * * @param chunkMetaData the corresponding chunk metadata - * @param timeFilter it should be a TimeFilter instead of a ValueFilter + * @param timeFilter it should be a TimeFilter instead of a ValueFilter */ public static List<IPageReader> loadPageReaderList(ChunkMetadata chunkMetaData, Filter timeFilter) throws IOException { diff --git a/server/src/main/java/org/apache/iotdb/db/utils/UpgradeUtils.java b/server/src/main/java/org/apache/iotdb/db/utils/UpgradeUtils.java index c213f64..1856157 100644 --- a/server/src/main/java/org/apache/iotdb/db/utils/UpgradeUtils.java +++ b/server/src/main/java/org/apache/iotdb/db/utils/UpgradeUtils.java @@ -68,7 +68,7 @@ public class UpgradeUtils { } try (TsFileSequenceReader tsFileSequenceReader = new TsFileSequenceReader( tsFileResource.getTsFile().getAbsolutePath())) { - if (tsFileSequenceReader.readVersionNumber().equals(TSFileConfig.VERSION_NUMBER_V1)) { + if (tsFileSequenceReader.readVersionNumber() == TSFileConfig.VERSION_NUMBER_V2.getBytes()[0]) { return true; } } catch (Exception e) { 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 76c1813..62267e6 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 @@ -223,4 +223,8 @@ public class ChunkHeader { this.dataSize += chunkHeader.getDataSize(); this.numOfPages += chunkHeader.getNumOfPages(); } + + public byte getChunkType() { + return chunkType; + } } 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 0b3e4cd..2c0acf9 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 @@ -62,6 +62,12 @@ public class PageHeader { return new PageHeader(uncompressedSize, compressedSize, statistics); } + public static PageHeader deserializeFrom(ByteBuffer buffer, Statistics chunkStatistic) { + int uncompressedSize = ReadWriteForEncodingUtils.readUnsignedVarInt(buffer); + int compressedSize = ReadWriteForEncodingUtils.readUnsignedVarInt(buffer); + return new PageHeader(uncompressedSize, compressedSize, chunkStatistic); + } + public int getUncompressedSize() { return uncompressedSize; } 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 a966ba4..5c7338d 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 @@ -226,12 +226,11 @@ public class TsFileSequenceReader implements AutoCloseable { /** * this function reads version number and checks compatibility of TsFile. */ - public String readVersionNumber() throws IOException { - ByteBuffer versionNumberBytes = ByteBuffer - .allocate(TSFileConfig.VERSION_NUMBER_V2.getBytes().length); - tsFileInput.read(versionNumberBytes, TSFileConfig.MAGIC_STRING.getBytes().length); - versionNumberBytes.flip(); - return new String(versionNumberBytes.array()); + public byte readVersionNumber() throws IOException { + ByteBuffer versionNumberByte = ByteBuffer.allocate(Byte.BYTES); + tsFileInput.read(versionNumberByte, TSFileConfig.MAGIC_STRING.getBytes().length); + versionNumberByte.flip(); + return versionNumberByte.get(); } /** @@ -728,7 +727,7 @@ public class TsFileSequenceReader implements AutoCloseable { ChunkHeader header = readChunkHeader(metaData.getOffsetOfChunkHeader(), chunkHeadSize); ByteBuffer buffer = readChunk(metaData.getOffsetOfChunkHeader() + header.getSerializedSize(), header.getDataSize()); - return new Chunk(header, buffer, metaData.getDeleteIntervalList()); + return new Chunk(header, buffer, metaData.getDeleteIntervalList(), metaData.getStatistics()); } /** @@ -911,8 +910,8 @@ public class TsFileSequenceReader implements AutoCloseable { if (fileSize < headerLength) { return TsFileCheckStatus.INCOMPATIBLE_FILE; } - if (!TSFileConfig.MAGIC_STRING.equals(readHeadMagic()) || !TSFileConfig.VERSION_NUMBER_V2 - .equals(readVersionNumber())) { + if (!TSFileConfig.MAGIC_STRING.equals(readHeadMagic()) || !(TSFileConfig.VERSION_NUMBER + == readVersionNumber())) { return TsFileCheckStatus.INCOMPATIBLE_FILE; } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/Chunk.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/Chunk.java index 1968aa9..a7bbadb 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/Chunk.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/Chunk.java @@ -23,6 +23,7 @@ import java.nio.ByteBuffer; import java.util.List; import org.apache.iotdb.tsfile.common.cache.Accountable; import org.apache.iotdb.tsfile.file.header.ChunkHeader; +import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics; /** * used in query. @@ -30,6 +31,7 @@ import org.apache.iotdb.tsfile.file.header.ChunkHeader; public class Chunk implements Accountable { private ChunkHeader chunkHeader; + private Statistics chunkStatistic; private ByteBuffer chunkData; /** * A list of deleted intervals. @@ -38,10 +40,11 @@ public class Chunk implements Accountable { private long ramSize; - public Chunk(ChunkHeader header, ByteBuffer buffer, List<TimeRange> deleteIntervalList) { + public Chunk(ChunkHeader header, ByteBuffer buffer, List<TimeRange> deleteIntervalList, Statistics chunkStatistic) { this.chunkHeader = header; this.chunkData = buffer; this.deleteIntervalList = deleteIntervalList; + this.chunkStatistic = chunkStatistic; } public ChunkHeader getHeader() { @@ -78,4 +81,8 @@ public class Chunk implements Accountable { public long getRamSize() { return ramSize; } + + public Statistics getChunkStatistic() { + return chunkStatistic; + } } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/controller/CachedChunkLoaderImpl.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/controller/CachedChunkLoaderImpl.java index 9c47e70..deb51b5 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/controller/CachedChunkLoaderImpl.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/controller/CachedChunkLoaderImpl.java @@ -61,7 +61,8 @@ public class CachedChunkLoaderImpl implements IChunkLoader { @Override public Chunk loadChunk(ChunkMetadata chunkMetaData) throws IOException { Chunk chunk = chunkCache.get(chunkMetaData); - return new Chunk(chunk.getHeader(), chunk.getData().duplicate(), chunk.getDeleteIntervalList()); + return new Chunk(chunk.getHeader(), chunk.getData().duplicate(), chunk.getDeleteIntervalList(), + chunkMetaData.getStatistics()); } @Override 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 f747926..729b8f7 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 @@ -19,26 +19,27 @@ package org.apache.iotdb.tsfile.read.reader.chunk; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.LinkedList; +import java.util.List; 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.ChunkHeader; import org.apache.iotdb.tsfile.file.header.PageHeader; 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.read.common.BatchData; import org.apache.iotdb.tsfile.read.common.Chunk; import org.apache.iotdb.tsfile.read.common.TimeRange; -import org.apache.iotdb.tsfile.read.reader.IPageReader; import org.apache.iotdb.tsfile.read.filter.basic.Filter; import org.apache.iotdb.tsfile.read.reader.IChunkReader; +import org.apache.iotdb.tsfile.read.reader.IPageReader; import org.apache.iotdb.tsfile.read.reader.page.PageReader; -import java.io.IOException; -import java.nio.ByteBuffer; -import java.util.LinkedList; -import java.util.List; - public class ChunkReader implements IChunkReader { private ChunkHeader chunkHeader; @@ -51,7 +52,7 @@ public class ChunkReader implements IChunkReader { protected Filter filter; private List<IPageReader> pageReaderList = new LinkedList<>(); - + private boolean isFromOldTsFile = false; /** @@ -72,11 +73,11 @@ public class ChunkReader implements IChunkReader { chunkHeader = chunk.getHeader(); this.unCompressor = IUnCompressor.getUnCompressor(chunkHeader.getCompressionType()); - - initAllPageReaders(); + initAllPageReaders(chunk.getChunkStatistic()); } - public ChunkReader(Chunk chunk, Filter filter, boolean isFromOldFile) throws IOException { + public ChunkReader(Chunk chunk, Filter filter, boolean isFromOldFile) + throws IOException { this.filter = filter; this.chunkDataBuffer = chunk.getData(); this.deleteIntervalList = chunk.getDeleteIntervalList(); @@ -84,14 +85,19 @@ public class ChunkReader implements IChunkReader { this.unCompressor = IUnCompressor.getUnCompressor(chunkHeader.getCompressionType()); this.isFromOldTsFile = isFromOldFile; - initAllPageReaders(); + initAllPageReaders(chunk.getChunkStatistic()); } - private void initAllPageReaders() throws IOException { + private void initAllPageReaders(Statistics chunkStatistic) throws IOException { // construct next satisfied page header while (chunkDataBuffer.remaining() > 0) { // deserialize a PageHeader from chunkDataBuffer - PageHeader pageHeader = PageHeader.deserializeFrom(chunkDataBuffer, chunkHeader.getDataType()); + PageHeader pageHeader; + if (chunkHeader.getChunkType() == MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER) { + pageHeader = PageHeader.deserializeFrom(chunkDataBuffer, chunkStatistic); + } else { + pageHeader = PageHeader.deserializeFrom(chunkDataBuffer, chunkHeader.getDataType()); + } // if the current page satisfies if (pageSatisfied(pageHeader)) { pageReaderList.add(constructPageReaderForNextPage(pageHeader)); @@ -102,7 +108,6 @@ public class ChunkReader implements IChunkReader { } - /** * judge if has next page whose page header satisfies the filter. */ @@ -156,9 +161,9 @@ public class ChunkReader implements IChunkReader { chunkDataBuffer.get(compressedPageBody); Decoder valueDecoder = Decoder - .getDecoderByType(chunkHeader.getEncodingType(), chunkHeader.getDataType()); + .getDecoderByType(chunkHeader.getEncodingType(), chunkHeader.getDataType()); byte[] uncompressedPageData = new byte[pageHeader.getUncompressedSize()]; - unCompressor.uncompress(compressedPageBody,0, compressedPageBodyLength, + unCompressor.uncompress(compressedPageBody, 0, compressedPageBodyLength, uncompressedPageData, 0); ByteBuffer pageData = ByteBuffer.wrap(uncompressedPageData); PageReader reader = new PageReader(pageHeader, pageData, chunkHeader.getDataType(), diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteIOUtils.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteIOUtils.java index 31eb8c7..ad251e5 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteIOUtils.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteIOUtils.java @@ -660,7 +660,7 @@ public class ReadWriteIOUtils { * string length's type is varInt */ public static String readVarIntString(ByteBuffer buffer) { - int strLength = readInt(buffer); + int strLength = ReadWriteForEncodingUtils.readVarInt(buffer); if (strLength < 0) { return null; } else if (strLength == 0) { diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java index 292e620..c7eefe0 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java @@ -48,7 +48,6 @@ import org.apache.iotdb.tsfile.read.common.Path; import org.apache.iotdb.tsfile.utils.BytesUtils; import org.apache.iotdb.tsfile.utils.Pair; import org.apache.iotdb.tsfile.utils.PublicBAOS; -import org.apache.iotdb.tsfile.utils.ReadWriteForEncodingUtils; import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; import org.apache.iotdb.tsfile.utils.VersionUtils; import org.apache.iotdb.tsfile.write.schema.MeasurementSchema; @@ -253,7 +252,7 @@ public class TsFileIOWriter { } // write TsFileMetaData size - ReadWriteForEncodingUtils.writeUnsignedVarInt(size, out.wrapAsStream());// write the size of the file metadata. + ReadWriteIOUtils.write(size, out.wrapAsStream());// write the size of the file metadata. // write magic string out.write(magicStringBytes);
