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 f51296c0bacdd711fdd519ae18db74aa26d62736 Author: JackieTien97 <[email protected]> AuthorDate: Wed Nov 25 15:46:44 2020 +0800 some changes --- .../iotdb/hadoop/tsfile/record/HDFSTSRecord.java | 4 +- .../org/apache/iotdb/db/metadata/MManager.java | 6 +-- .../db/qp/physical/crud/InsertTabletPlan.java | 4 +- .../iotdb/tsfile/file/metadata/ChunkMetadata.java | 21 +++++----- .../tsfile/file/metadata/MetadataIndexEntry.java | 4 +- .../tsfile/file/metadata/MetadataIndexNode.java | 5 ++- .../tsfile/file/metadata/TimeseriesMetadata.java | 9 +++-- .../iotdb/tsfile/file/metadata/TsFileMetadata.java | 20 +++++----- .../file/metadata/enums/CompressionType.java | 41 +++----------------- .../tsfile/file/metadata/enums/TSDataType.java | 45 +++------------------- .../tsfile/file/metadata/enums/TSEncoding.java | 40 +++++++++++-------- .../iotdb/tsfile/read/TsFileSequenceReader.java | 21 +++++++--- .../iotdb/tsfile/utils/ReadWriteIOUtils.java | 31 +++++++-------- .../tsfile/write/schema/MeasurementSchema.java | 24 ++++++------ .../iotdb/tsfile/write/writer/TsFileIOWriter.java | 14 ------- 15 files changed, 112 insertions(+), 177 deletions(-) diff --git a/hadoop/src/main/java/org/apache/iotdb/hadoop/tsfile/record/HDFSTSRecord.java b/hadoop/src/main/java/org/apache/iotdb/hadoop/tsfile/record/HDFSTSRecord.java index f996094..a5d380e 100644 --- a/hadoop/src/main/java/org/apache/iotdb/hadoop/tsfile/record/HDFSTSRecord.java +++ b/hadoop/src/main/java/org/apache/iotdb/hadoop/tsfile/record/HDFSTSRecord.java @@ -100,7 +100,7 @@ public class HDFSTSRecord implements Writable { out.write(deviceId.getBytes(StandardCharsets.UTF_8)); out.writeInt(dataPointList.size()); for (DataPoint dataPoint : dataPointList) { - out.writeShort(dataPoint.getType().serialize()); + out.write(dataPoint.getType().serialize()); out.writeInt(dataPoint.getMeasurementId().getBytes(StandardCharsets.UTF_8).length); out.write(dataPoint.getMeasurementId().getBytes(StandardCharsets.UTF_8)); switch (dataPoint.getType()) { @@ -139,7 +139,7 @@ public class HDFSTSRecord implements Writable { List<DataPoint> dataPoints = new ArrayList<>(len); for (int i = 0; i < len; i++) { - TSDataType dataType = TSDataType.deserialize(in.readShort()); + TSDataType dataType = TSDataType.deserialize(in.readByte()); int lenOfMeasurementId = in.readInt(); byte[] c = new byte[lenOfMeasurementId]; in.readFully(c); diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java index d349f21..488b135 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java @@ -350,9 +350,9 @@ public class MManager { } CreateTimeSeriesPlan plan = new CreateTimeSeriesPlan(new PartialPath(args[1]), - TSDataType.deserialize(Short.parseShort(args[2])), - TSEncoding.deserialize(Short.parseShort(args[3])), - CompressionType.deserialize(Short.parseShort(args[4])), props, tagMap, null, alias); + TSDataType.deserialize(Byte.parseByte(args[2])), + TSEncoding.deserialize(Byte.parseByte(args[3])), + CompressionType.deserialize(Byte.parseByte(args[4])), props, tagMap, null, alias); createTimeseries(plan, offset); break; diff --git a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java index daa9d84..abdbe6f 100644 --- a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java +++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java @@ -165,7 +165,7 @@ public class InsertTabletPlan extends InsertPlan { if (dataType == null) { continue; } - stream.writeShort(dataType.serialize()); + stream.write(dataType.serialize()); } } @@ -403,7 +403,7 @@ public class InsertTabletPlan extends InsertPlan { this.dataTypes = new TSDataType[measurementSize]; for (int i = 0; i < measurementSize; i++) { - dataTypes[i] = TSDataType.deserialize(buffer.getShort()); + dataTypes[i] = TSDataType.deserialize(buffer.get()); } int rows = buffer.getInt(); diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java index 71e2596..6f494b3 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java @@ -83,9 +83,9 @@ public class ChunkMetadata implements Accountable { * constructor of ChunkMetaData. * * @param measurementUid measurement id - * @param tsDataType time series data type - * @param fileOffset file offset - * @param statistics value statistics + * @param tsDataType time series data type + * @param fileOffset file offset + * @param statistics value statistics */ public ChunkMetadata(String measurementUid, TSDataType tsDataType, long fileOffset, Statistics statistics) { @@ -144,10 +144,7 @@ public class ChunkMetadata implements Accountable { */ public int serializeTo(OutputStream outputStream) throws IOException { int byteLen = 0; - - byteLen += ReadWriteIOUtils.write(measurementUid, outputStream); byteLen += ReadWriteIOUtils.write(offsetOfChunkHeader, outputStream); - byteLen += ReadWriteIOUtils.write(tsDataType, outputStream); byteLen += statistics.serialize(outputStream); return byteLen; } @@ -155,16 +152,18 @@ public class ChunkMetadata implements Accountable { /** * deserialize from ByteBuffer. * - * @param buffer ByteBuffer + * @param buffer ByteBuffer + * @param measurementUid measurementUid of this chunk metadata + * @param tsDataType data type of this chunk metadata * @return ChunkMetaData object */ - public static ChunkMetadata deserializeFrom(ByteBuffer buffer) { + public static ChunkMetadata deserializeFrom(ByteBuffer buffer, String measurementUid, + TSDataType tsDataType) { ChunkMetadata chunkMetaData = new ChunkMetadata(); - chunkMetaData.measurementUid = ReadWriteIOUtils.readString(buffer); + chunkMetaData.measurementUid = measurementUid; + chunkMetaData.tsDataType = tsDataType; chunkMetaData.offsetOfChunkHeader = ReadWriteIOUtils.readLong(buffer); - chunkMetaData.tsDataType = ReadWriteIOUtils.readDataType(buffer); - chunkMetaData.statistics = Statistics.deserialize(buffer, chunkMetaData.tsDataType); return chunkMetaData; diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexEntry.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexEntry.java index 5325992..73667ba 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexEntry.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexEntry.java @@ -56,13 +56,13 @@ public class MetadataIndexEntry { public int serializeTo(OutputStream outputStream) throws IOException { int byteLen = 0; - byteLen += ReadWriteIOUtils.write(name, outputStream); + byteLen += ReadWriteIOUtils.writeVar(name, outputStream); byteLen += ReadWriteIOUtils.write(offset, outputStream); return byteLen; } public static MetadataIndexEntry deserializeFrom(ByteBuffer buffer) { - String name = ReadWriteIOUtils.readString(buffer); + String name = ReadWriteIOUtils.readVarIntString(buffer); long offset = ReadWriteIOUtils.readLong(buffer); return new MetadataIndexEntry(name, offset); } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexNode.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexNode.java index 95f67f5..0600fd4 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexNode.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexNode.java @@ -27,6 +27,7 @@ import java.util.List; import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor; import org.apache.iotdb.tsfile.file.metadata.enums.MetadataIndexNodeType; import org.apache.iotdb.tsfile.utils.Pair; +import org.apache.iotdb.tsfile.utils.ReadWriteForEncodingUtils; import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; public class MetadataIndexNode { @@ -87,7 +88,7 @@ public class MetadataIndexNode { public int serializeTo(OutputStream outputStream) throws IOException { int byteLen = 0; - byteLen += ReadWriteIOUtils.write(children.size(), outputStream); + byteLen += ReadWriteForEncodingUtils.writeUnsignedVarInt(children.size(), outputStream); for (MetadataIndexEntry metadataIndexEntry : children) { byteLen += metadataIndexEntry.serializeTo(outputStream); } @@ -98,7 +99,7 @@ public class MetadataIndexNode { public static MetadataIndexNode deserializeFrom(ByteBuffer buffer) { List<MetadataIndexEntry> children = new ArrayList<>(); - int size = ReadWriteIOUtils.readInt(buffer); + int size = ReadWriteForEncodingUtils.readUnsignedVarInt(buffer); for (int i = 0; i < size; i++) { children.add(MetadataIndexEntry.deserializeFrom(buffer)); } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java index 0869643..6188e1d 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java @@ -27,6 +27,7 @@ import org.apache.iotdb.tsfile.common.cache.Accountable; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics; import org.apache.iotdb.tsfile.read.controller.IChunkMetadataLoader; +import org.apache.iotdb.tsfile.utils.ReadWriteForEncodingUtils; import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; public class TimeseriesMetadata implements Accountable { @@ -75,7 +76,8 @@ public class TimeseriesMetadata implements Accountable { timeseriesMetaData.setMeasurementId(ReadWriteIOUtils.readString(buffer)); timeseriesMetaData.setTSDataType(ReadWriteIOUtils.readDataType(buffer)); timeseriesMetaData.setOffsetOfChunkMetaDataList(ReadWriteIOUtils.readLong(buffer)); - timeseriesMetaData.setDataSizeOfChunkMetaDataList(ReadWriteIOUtils.readInt(buffer)); + timeseriesMetaData + .setDataSizeOfChunkMetaDataList(ReadWriteForEncodingUtils.readUnsignedVarInt(buffer)); timeseriesMetaData.setStatistics(Statistics.deserialize(buffer, timeseriesMetaData.dataType)); return timeseriesMetaData; } @@ -92,7 +94,8 @@ public class TimeseriesMetadata implements Accountable { byteLen += ReadWriteIOUtils.write(measurementId, outputStream); byteLen += ReadWriteIOUtils.write(dataType, outputStream); byteLen += ReadWriteIOUtils.write(startOffsetOfChunkMetaDataList, outputStream); - byteLen += ReadWriteIOUtils.write(chunkMetaDataListDataSize, outputStream); + byteLen += ReadWriteForEncodingUtils + .writeUnsignedVarInt(chunkMetaDataListDataSize, outputStream); byteLen += statistics.serialize(outputStream); return byteLen; } @@ -161,7 +164,7 @@ public class TimeseriesMetadata implements Accountable { public long getRamSize() { return ramSize; } - + public void setSeq(boolean seq) { isSeq = seq; } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TsFileMetadata.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TsFileMetadata.java index 33c5f8d..2831aa5 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TsFileMetadata.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TsFileMetadata.java @@ -77,9 +77,9 @@ public class TsFileMetadata { // read bloom filter if (buffer.hasRemaining()) { - byte[] bytes = ReadWriteIOUtils.readByteBufferWithSelfDescriptionLength(buffer).array(); - int filterSize = ReadWriteIOUtils.readInt(buffer); - int hashFunctionSize = ReadWriteIOUtils.readInt(buffer); + byte[] bytes = ReadWriteIOUtils.readByteBufferWithSelfDescriptionLength(buffer); + int filterSize = ReadWriteForEncodingUtils.readUnsignedVarInt(buffer); + int hashFunctionSize = ReadWriteForEncodingUtils.readUnsignedVarInt(buffer); fileMetaData.bloomFilter = BloomFilter.buildBloomFilter(bytes, filterSize, hashFunctionSize); } @@ -103,15 +103,12 @@ public class TsFileMetadata { if (metadataIndex != null) { byteLen += metadataIndex.serializeTo(outputStream); } else { + // TODO Maybe should throw an exception byteLen += ReadWriteIOUtils.write(0, outputStream); } - // totalChunkNum, invalidChunkNum - byteLen += ReadWriteIOUtils.write(totalChunkNum, outputStream); - byteLen += ReadWriteIOUtils.write(invalidChunkNum, outputStream); - // versionInfo - byteLen += ReadWriteIOUtils.write(versionInfo.size(), outputStream); + byteLen += ReadWriteForEncodingUtils.writeUnsignedVarInt(versionInfo.size(), outputStream); for (Pair<Long, Long> versionPair : versionInfo) { byteLen += ReadWriteIOUtils.write(versionPair.left, outputStream); byteLen += ReadWriteIOUtils.write(versionPair.right, outputStream); @@ -135,11 +132,12 @@ public class TsFileMetadata { BloomFilter filter = buildBloomFilter(paths); byte[] bytes = filter.serialize(); - byteLen += ReadWriteIOUtils.write(bytes.length, outputStream); + byteLen += ReadWriteForEncodingUtils.writeUnsignedVarInt(bytes.length, outputStream); outputStream.write(bytes); byteLen += bytes.length; - byteLen += ReadWriteIOUtils.write(filter.getSize(), outputStream); - byteLen += ReadWriteIOUtils.write(filter.getHashFunctionSize(), outputStream); + byteLen += ReadWriteForEncodingUtils.writeUnsignedVarInt(filter.getSize(), outputStream); + byteLen += ReadWriteForEncodingUtils + .writeUnsignedVarInt(filter.getHashFunctionSize(), outputStream); return byteLen; } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/CompressionType.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/CompressionType.java index e8380e6..5976e09 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/CompressionType.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/CompressionType.java @@ -22,24 +22,12 @@ public enum CompressionType { UNCOMPRESSED, SNAPPY, GZIP, LZO, SDT, PAA, PLA, LZ4; /** - * deserialize short number. + * deserialize byte number. * - * @param compressor short number + * @param compressor byte number * @return CompressionType */ - public static CompressionType deserialize(short compressor) { - return getCompressionType(compressor); - } - - public static byte deserializeToByte(short compressor) { - if (compressor >= 8 || compressor < 0) { - throw new IllegalArgumentException("Invalid input: " + compressor); - } - return (byte) compressor; - } - - - private static CompressionType getCompressionType(short compressor) { + public static CompressionType deserialize(byte compressor) { if (compressor >= 8 || compressor < 0) { throw new IllegalArgumentException("Invalid input: " + compressor); } @@ -63,33 +51,14 @@ public enum CompressionType { } } - /** - * give an byte to return a compression type. - * - * @param compressor byte number - * @return CompressionType - */ - public static CompressionType byteToEnum(byte compressor) { - return getCompressionType(compressor); - } - public static int getSerializedSize() { - return Short.BYTES; - } - - /** - * serialize. - * - * @return short number - */ - public short serialize() { - return enumToByte(); + return Byte.BYTES; } /** * @return byte number */ - public byte enumToByte() { + public byte serialize() { switch (this) { case SNAPPY: return 1; diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSDataType.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSDataType.java index 505be1b..79dd1e0 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSDataType.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSDataType.java @@ -32,12 +32,7 @@ public enum TSDataType { * @param type -param to judge enum type * @return -enum type */ - public static TSDataType deserialize(short type) { - return getTsDataType(type); - } - - - private static TSDataType getTsDataType(short type) { + public static TSDataType deserialize(byte type) { if (type >= 6 || type < 0) { throw new IllegalArgumentException("Invalid input: " + type); } @@ -57,46 +52,16 @@ public enum TSDataType { } } - public static byte deserializeToByte(short type) { - if (type >= 6 || type < 0) { - throw new IllegalArgumentException("Invalid input: " + type); - } - return (byte) type; - } - - /** - * give an byte to return a data type. - * - * @param type byte number - * @return data type - */ - public static TSDataType byteToEnum(byte type) { - return getTsDataType(type); - } - - public static TSDataType deserializeFrom(ByteBuffer buffer) { - return deserialize(buffer.getShort()); - } - public static int getSerializedSize() { - return Short.BYTES; + return Byte.BYTES; } public void serializeTo(ByteBuffer byteBuffer) { - byteBuffer.putShort(serialize()); + byteBuffer.put(serialize()); } public void serializeTo(DataOutputStream outputStream) throws IOException { - outputStream.writeShort(serialize()); - } - - /** - * return a serialize data type. - * - * @return -enum type - */ - public short serialize() { - return enumToByte(); + outputStream.write(serialize()); } public int getDataTypeSize() { @@ -119,7 +84,7 @@ public enum TSDataType { /** * @return byte number */ - public byte enumToByte() { + public byte serialize() { switch (this) { case BOOLEAN: return 0; diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSEncoding.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSEncoding.java index a909e47..9232fb9 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSEncoding.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/enums/TSEncoding.java @@ -28,18 +28,33 @@ public enum TSEncoding { * @param encoding -use to determine encoding type * @return -encoding type */ - public static TSEncoding deserialize(short encoding) { - return getTsEncoding(encoding); - } - - public static byte deserializeToByte(short encoding) { + public static TSEncoding deserialize(byte encoding) { if (encoding < 0 || 8 < encoding) { throw new IllegalArgumentException("Invalid input: " + encoding); } - return (byte) encoding; + switch (encoding) { + case 1: + return PLAIN_DICTIONARY; + case 2: + return RLE; + case 3: + return DIFF; + case 4: + return TS_2DIFF; + case 5: + return BITMAP; + case 6: + return GORILLA_V1; + case 7: + return REGULAR; + case 8: + return GORILLA; + default: + return PLAIN; + } } - private static TSEncoding getTsEncoding(short encoding) { + private static TSEncoding getTsEncoding(byte encoding) { if (encoding < 0 || 8 < encoding) { throw new IllegalArgumentException("Invalid input: " + encoding); } @@ -76,7 +91,7 @@ public enum TSEncoding { } public static int getSerializedSize() { - return Short.BYTES; + return Byte.BYTES; } /** @@ -84,14 +99,7 @@ public enum TSEncoding { * * @return -encoding type */ - public short serialize() { - return enumToByte(); - } - - /** - * @return byte number - */ - public byte enumToByte() { + public byte serialize() { switch (this) { case PLAIN_DICTIONARY: return 1; 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 9b9b3f4..16e2b71 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 @@ -491,9 +491,7 @@ public class TsFileSequenceReader implements AutoCloseable { */ public Map<String, List<ChunkMetadata>> readChunkMetadataInDevice(String device) throws IOException { - if (tsFileMetaData == null) { - readFileMetadata(); - } + readFileMetadata(); long start = 0; int size = 0; @@ -507,8 +505,16 @@ public class TsFileSequenceReader implements AutoCloseable { // read buffer of all ChunkMetadatas of this device ByteBuffer buffer = readData(start, size); Map<String, List<ChunkMetadata>> seriesMetadata = new HashMap<>(); + int index = 0; + int curSize = timeseriesMetadataMap.get(index).getDataSizeOfChunkMetaDataList(); while (buffer.hasRemaining()) { - ChunkMetadata chunkMetadata = ChunkMetadata.deserializeFrom(buffer); + if (buffer.position() >= curSize) { + index++; + curSize += timeseriesMetadataMap.get(index).getDataSizeOfChunkMetaDataList(); + } + ChunkMetadata chunkMetadata = ChunkMetadata + .deserializeFrom(buffer, timeseriesMetadataMap.get(index).getMeasurementId(), + timeseriesMetadataMap.get(index).getTSDataType()); seriesMetadata.computeIfAbsent(chunkMetadata.getMeasurementUid(), key -> new ArrayList<>()) .add(chunkMetadata); } @@ -693,7 +699,8 @@ public class TsFileSequenceReader implements AutoCloseable { /** * read the chunk's header. - * @param position the file offset of this chunk's header + * + * @param position the file offset of this chunk's header * @param chunkHeaderSize the size of chunk's header */ private ChunkHeader readChunkHeader(long position, int chunkHeaderSize) throws IOException { @@ -1034,7 +1041,9 @@ public class TsFileSequenceReader implements AutoCloseable { ByteBuffer buffer = readData(startOffsetOfChunkMetadataList, dataSizeOfChunkMetadataList); while (buffer.hasRemaining()) { - chunkMetadataList.add(ChunkMetadata.deserializeFrom(buffer)); + chunkMetadataList.add(ChunkMetadata + .deserializeFrom(buffer, timeseriesMetaData.getMeasurementId(), + timeseriesMetaData.getTSDataType())); } VersionUtils.applyVersion(chunkMetadataList, versionInfo); 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 524ee01..a5b02de 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 @@ -446,7 +446,7 @@ public class ReadWriteIOUtils { */ public static int write(CompressionType compressionType, OutputStream outputStream) throws IOException { - short n = compressionType.serialize(); + byte n = compressionType.serialize(); return write(n, outputStream); } @@ -454,7 +454,7 @@ public class ReadWriteIOUtils { * write compressionType to byteBuffer. */ public static int write(CompressionType compressionType, ByteBuffer buffer) { - short n = compressionType.serialize(); + byte n = compressionType.serialize(); return write(n, buffer); } @@ -462,12 +462,12 @@ public class ReadWriteIOUtils { * TSDataType. */ public static int write(TSDataType dataType, OutputStream outputStream) throws IOException { - short n = dataType.serialize(); + byte n = dataType.serialize(); return write(n, outputStream); } public static int write(TSDataType dataType, ByteBuffer buffer) { - short n = dataType.serialize(); + byte n = dataType.serialize(); return write(n, buffer); } @@ -475,12 +475,12 @@ public class ReadWriteIOUtils { * TSEncoding. */ public static int write(TSEncoding encoding, OutputStream outputStream) throws IOException { - short n = encoding.serialize(); + byte n = encoding.serialize(); return write(n, outputStream); } public static int write(TSEncoding encoding, ByteBuffer buffer) { - short n = encoding.serialize(); + byte n = encoding.serialize(); return write(n, buffer); } @@ -768,14 +768,11 @@ public class ReadWriteIOUtils { * <p> * read a int + buffer */ - public static ByteBuffer readByteBufferWithSelfDescriptionLength(ByteBuffer buffer) { + public static byte[] readByteBufferWithSelfDescriptionLength(ByteBuffer buffer) { int byteLength = ReadWriteForEncodingUtils.readUnsignedVarInt(buffer); byte[] bytes = new byte[byteLength]; buffer.get(bytes); - ByteBuffer byteBuffer = ByteBuffer.allocate(byteLength); - byteBuffer.put(bytes); - byteBuffer.flip(); - return byteBuffer; + return bytes; } /** @@ -913,32 +910,32 @@ public class ReadWriteIOUtils { public static CompressionType readCompressionType(InputStream inputStream) throws IOException { byte n = readByte(inputStream); - return CompressionType.byteToEnum(n); + return CompressionType.deserialize(n); } public static CompressionType readCompressionType(ByteBuffer buffer) { byte n = buffer.get(); - return CompressionType.byteToEnum(n); + return CompressionType.deserialize(n); } public static TSDataType readDataType(InputStream inputStream) throws IOException { byte n = readByte(inputStream); - return TSDataType.byteToEnum(n); + return TSDataType.deserialize(n); } public static TSDataType readDataType(ByteBuffer buffer) { byte n = buffer.get(); - return TSDataType.byteToEnum(n); + return TSDataType.deserialize(n); } public static TSEncoding readEncoding(InputStream inputStream) throws IOException { byte n = readByte(inputStream); - return TSEncoding.byteToEnum(n); + return TSEncoding.deserialize(n); } public static TSEncoding readEncoding(ByteBuffer buffer) { byte n = buffer.get(); - return TSEncoding.byteToEnum(n); + return TSEncoding.deserialize(n); } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/schema/MeasurementSchema.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/schema/MeasurementSchema.java index 11030c8..4ec4a23 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/schema/MeasurementSchema.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/schema/MeasurementSchema.java @@ -87,11 +87,11 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali */ public MeasurementSchema(String measurementId, TSDataType type, TSEncoding encoding, CompressionType compressionType, Map<String, String> props) { - this.type = type.enumToByte(); + this.type = type.serialize(); this.measurementId = measurementId; - this.encoding = encoding.enumToByte(); + this.encoding = encoding.serialize(); this.props = props; - this.compressor = compressionType.enumToByte(); + this.compressor = compressionType.serialize(); } public MeasurementSchema(String measurementId, byte type, byte encoding, @@ -111,11 +111,11 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali measurementSchema.measurementId = ReadWriteIOUtils.readString(inputStream); - measurementSchema.type = TSDataType.deserializeToByte(ReadWriteIOUtils.readShort(inputStream)); + measurementSchema.type = ReadWriteIOUtils.readByte(inputStream); - measurementSchema.encoding = TSEncoding.deserializeToByte(ReadWriteIOUtils.readShort(inputStream)); + measurementSchema.encoding = ReadWriteIOUtils.readByte(inputStream); - measurementSchema.compressor = CompressionType.deserializeToByte(ReadWriteIOUtils.readShort(inputStream)); + measurementSchema.compressor = ReadWriteIOUtils.readByte(inputStream); int size = ReadWriteIOUtils.readInt(inputStream); if (size > 0) { @@ -140,11 +140,11 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali measurementSchema.measurementId = ReadWriteIOUtils.readString(buffer); - measurementSchema.type = TSDataType.deserializeToByte(ReadWriteIOUtils.readShort(buffer)); + measurementSchema.type = ReadWriteIOUtils.readByte(buffer); - measurementSchema.encoding = TSEncoding.deserializeToByte(ReadWriteIOUtils.readShort(buffer)); + measurementSchema.encoding = ReadWriteIOUtils.readByte(buffer); - measurementSchema.compressor = CompressionType.deserializeToByte(ReadWriteIOUtils.readShort(buffer)); + measurementSchema.compressor = ReadWriteIOUtils.readByte(buffer); int size = ReadWriteIOUtils.readInt(buffer); if (size > 0) { @@ -178,7 +178,7 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali } public TSDataType getType() { - return TSDataType.byteToEnum(type); + return TSDataType.deserialize(type); } public void setProps(Map<String, String> props) { @@ -212,7 +212,7 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali } public CompressionType getCompressor() { - return CompressionType.byteToEnum(compressor); + return CompressionType.deserialize(compressor); } /** @@ -312,6 +312,6 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali } public void setType(TSDataType type) { - this.type = (byte) type.serialize(); + this.type = type.serialize(); } } 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 40a9161..3c83308 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 @@ -72,8 +72,6 @@ public class TsFileIOWriter { protected TsFileOutput out; protected boolean canWrite = true; - protected int totalChunkNum = 0; - protected int invalidChunkNum; protected File file; // current flushed Chunk @@ -207,7 +205,6 @@ public class TsFileIOWriter { public void endCurrentChunk() { chunkMetadataList.add(currentChunkMetadata); currentChunkMetadata = null; - totalChunkNum++; } /** @@ -234,8 +231,6 @@ public class TsFileIOWriter { TsFileMetadata tsFileMetaData = new TsFileMetadata(); tsFileMetaData.setMetadataIndex(metadataIndex); tsFileMetaData.setVersionInfo(versionInfo); - tsFileMetaData.setTotalChunkNum(totalChunkNum); - tsFileMetaData.setInvalidChunkNum(invalidChunkNum); tsFileMetaData.setMetaOffset(metaOffset); long footerIndex = out.getPosition(); @@ -359,14 +354,6 @@ public class TsFileIOWriter { out.write(new byte[]{MetaMarker.CHUNK_HEADER}); } - public int getTotalChunkNum() { - return totalChunkNum; - } - - public int getInvalidChunkNum() { - return invalidChunkNum; - } - public File getFile() { return file; } @@ -400,7 +387,6 @@ public class TsFileIOWriter { if (!chunkValid) { chunkMetaDataIterator.remove(); chunkNum--; - invalidChunkNum++; } else { startTimeIdxes.put(path, startTimeIdx + 1); }
