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 0b153b20404515f69889bb77fde579f5804dd6e4 Author: JackieTien97 <[email protected]> AuthorDate: Wed Dec 2 16:03:52 2020 +0800 fix bug --- .../file/metadata/MetadataIndexConstructor.java | 8 +++--- .../tsfile/file/metadata/MetadataIndexNode.java | 6 ++--- .../iotdb/tsfile/read/TsFileSequenceReader.java | 29 ++++++++++++++++++---- .../tsfile/utils/ReadWriteForEncodingUtils.java | 5 ++-- .../iotdb/tsfile/utils/ReadWriteIOUtils.java | 4 +-- .../tsfile/write/schema/MeasurementSchema.java | 12 ++++----- .../iotdb/tsfile/read/GetAllDevicesTest.java | 14 +++-------- .../tsfile/read/TsFileSequenceReaderTest.java | 13 ++++++---- .../iotdb/tsfile/write/TsFileIOWriterTest.java | 7 +++--- .../iotdb/tsfile/write/writer/PageWriterTest.java | 4 +-- .../write/writer/RestorableTsFileIOWriterTest.java | 5 ++-- 11 files changed, 63 insertions(+), 44 deletions(-) diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexConstructor.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexConstructor.java index d3070be..df0aa4d 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexConstructor.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/MetadataIndexConstructor.java @@ -26,14 +26,14 @@ import java.util.Map; import java.util.Map.Entry; import java.util.Queue; import java.util.TreeMap; +import org.apache.iotdb.tsfile.common.conf.TSFileConfig; import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor; import org.apache.iotdb.tsfile.file.metadata.enums.MetadataIndexNodeType; import org.apache.iotdb.tsfile.write.writer.TsFileOutput; public class MetadataIndexConstructor { - private static final int MAX_DEGREE_OF_INDEX_NODE = TSFileDescriptor.getInstance().getConfig() - .getMaxDegreeOfIndexNode(); + private static final TSFileConfig config = TSFileDescriptor.getInstance().getConfig(); private MetadataIndexConstructor() { throw new IllegalStateException("Utility class"); @@ -62,7 +62,7 @@ public class MetadataIndexConstructor { for (int i = 0; i < entry.getValue().size(); i++) { timeseriesMetadata = entry.getValue().get(i); // when constructing from leaf node, every "degree number of nodes" are related to an entry - if (i % MAX_DEGREE_OF_INDEX_NODE == 0) { + if (i % config.getMaxDegreeOfIndexNode() == 0) { if (currentIndexNode.isFull()) { addCurrentIndexNodeToQueue(currentIndexNode, measurementMetadataIndexQueue, out); currentIndexNode = new MetadataIndexNode(MetadataIndexNodeType.LEAF_MEASUREMENT); @@ -78,7 +78,7 @@ public class MetadataIndexConstructor { } // if not exceed the max child nodes num, ignore the device index and directly point to the measurement - if (deviceMetadataIndexMap.size() <= MAX_DEGREE_OF_INDEX_NODE) { + if (deviceMetadataIndexMap.size() <= config.getMaxDegreeOfIndexNode()) { MetadataIndexNode metadataIndexNode = new MetadataIndexNode( MetadataIndexNodeType.INTERNAL_MEASUREMENT); for (Map.Entry<String, MetadataIndexNode> entry : deviceMetadataIndexMap.entrySet()) { 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 0600fd4..f3ff04f 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 @@ -24,6 +24,7 @@ import java.io.OutputStream; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; +import org.apache.iotdb.tsfile.common.conf.TSFileConfig; import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor; import org.apache.iotdb.tsfile.file.metadata.enums.MetadataIndexNodeType; import org.apache.iotdb.tsfile.utils.Pair; @@ -32,8 +33,7 @@ import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; public class MetadataIndexNode { - private static final int MAX_DEGREE_OF_INDEX_NODE = TSFileDescriptor.getInstance().getConfig() - .getMaxDegreeOfIndexNode(); + private static final TSFileConfig config = TSFileDescriptor.getInstance().getConfig(); private List<MetadataIndexEntry> children; private long endOffset; @@ -76,7 +76,7 @@ public class MetadataIndexNode { } boolean isFull() { - return children.size() == MAX_DEGREE_OF_INDEX_NODE; + return children.size() == config.getMaxDegreeOfIndexNode(); } MetadataIndexEntry peek() { 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 c3d2c5a..d15a194 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 @@ -209,8 +209,14 @@ 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 + Byte.BYTES - && (readTailMagic().equals(readHeadMagic())); + long size = tsFileInput.size(); + if (size >= TSFileConfig.MAGIC_STRING.getBytes().length * 2 + Byte.BYTES) { + String tailMagic = readTailMagic(); + String headMagic = readHeadMagic(); + return tailMagic.equals(headMagic); + } else { + return false; + } } /** @@ -999,6 +1005,9 @@ public class TsFileSequenceReader implements AutoCloseable { chunkMetadataList.add(currentChunk); break; case MetaMarker.CHUNK_GROUP_HEADER: + // if there is something wrong with the ChunkGroup Header, we will drop this ChunkGroup + // because we can not guarantee the correctness of the deviceId. + truncatedSize = this.position() - 1; if (lastDeviceId != null) { // schema of last chunk group if (newSchema != null) { @@ -1011,15 +1020,13 @@ public class TsFileSequenceReader implements AutoCloseable { // 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(); lastDeviceId = chunkGroupHeader.getDeviceID(); break; case MetaMarker.VERSION: + truncatedSize = this.position() - 1; if (lastDeviceId != null) { // schema of last chunk group if (newSchema != null) { @@ -1046,6 +1053,18 @@ public class TsFileSequenceReader implements AutoCloseable { } // now we read the tail of the data section, so we are sure that the last // ChunkGroupFooter is complete. + 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)); + } truncatedSize = this.position() - 1; } catch (Exception e) { logger.info("TsFile {} self-check cannot proceed at position {} " + "recovered, because : {}", diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteForEncodingUtils.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteForEncodingUtils.java index 578ab4c..d3783df 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteForEncodingUtils.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/utils/ReadWriteForEncodingUtils.java @@ -98,10 +98,11 @@ public class ReadWriteForEncodingUtils { public static int readUnsignedVarInt(InputStream in) throws IOException { int value = 0; int i = 0; - int b; - while (((b = in.read()) & 0x80) != 0) { + int b = in.read(); + while (b != -1 && (b & 0x80) != 0) { value |= (b & 0x7F) << i; i += 7; + b = in.read(); } return value | (b << i); } 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 ad251e5..17c0935 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 @@ -390,11 +390,11 @@ public class ReadWriteIOUtils { public static int writeVar(String s, ByteBuffer buffer) { if (s == null) { - return write(-1, buffer); + return ReadWriteForEncodingUtils.writeVarInt(-1, buffer); } int len = 0; byte[] bytes = s.getBytes(); - len += write(bytes.length, buffer); + len += ReadWriteForEncodingUtils.writeVarInt(bytes.length, buffer); buffer.put(bytes); len += bytes.length; return len; 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 18127d2..7a6ca04 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 @@ -223,11 +223,11 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali byteLen += ReadWriteIOUtils.write(measurementId, outputStream); - byteLen += ReadWriteIOUtils.write((short) type, outputStream); + byteLen += ReadWriteIOUtils.write(type, outputStream); - byteLen += ReadWriteIOUtils.write((short) encoding, outputStream); + byteLen += ReadWriteIOUtils.write(encoding, outputStream); - byteLen += ReadWriteIOUtils.write((short) compressor, outputStream); + byteLen += ReadWriteIOUtils.write(compressor, outputStream); if (props == null) { byteLen += ReadWriteIOUtils.write(0, outputStream); @@ -250,11 +250,11 @@ public class MeasurementSchema implements Comparable<MeasurementSchema>, Seriali byteLen += ReadWriteIOUtils.write(measurementId, buffer); - byteLen += ReadWriteIOUtils.write((short) type, buffer); + byteLen += ReadWriteIOUtils.write(type, buffer); - byteLen += ReadWriteIOUtils.write((short) encoding, buffer); + byteLen += ReadWriteIOUtils.write(encoding, buffer); - byteLen += ReadWriteIOUtils.write((short) compressor, buffer); + byteLen += ReadWriteIOUtils.write(compressor, buffer); if (props == null) { byteLen += ReadWriteIOUtils.write(0, buffer); diff --git a/tsfile/src/test/java/org/apache/iotdb/tsfile/read/GetAllDevicesTest.java b/tsfile/src/test/java/org/apache/iotdb/tsfile/read/GetAllDevicesTest.java index 05add5c..61c6465 100644 --- a/tsfile/src/test/java/org/apache/iotdb/tsfile/read/GetAllDevicesTest.java +++ b/tsfile/src/test/java/org/apache/iotdb/tsfile/read/GetAllDevicesTest.java @@ -42,7 +42,7 @@ public class GetAllDevicesTest { } @After - public void after() throws IOException { + public void after() { FileGenerator.after(); conf.setMaxDegreeOfIndexNode(maxDegreeOfIndexNode); } @@ -70,19 +70,13 @@ public class GetAllDevicesTest { public void testGetAllDevices(int deviceNum, int measurementNum) throws IOException { FileGenerator.generateFile(10000, deviceNum, measurementNum); try (TsFileSequenceReader fileReader = new TsFileSequenceReader(FILE_PATH)) { - ReadOnlyTsFile tsFile = new ReadOnlyTsFile(fileReader); - - // test - try (TsFileSequenceReader reader = new TsFileSequenceReader(FILE_PATH)) { - List<String> devices = reader.getAllDevices(); + + List<String> devices = fileReader.getAllDevices(); Assert.assertEquals(deviceNum, devices.size()); for (int i = 0; i < deviceNum; i++) { Assert.assertTrue(devices.contains("d" + i)); } - } - - // after - tsFile.close(); + FileGenerator.after(); } } 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 4e8b7e0..1e67454 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 @@ -20,6 +20,7 @@ package org.apache.iotdb.tsfile.read; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.HashMap; import java.util.HashSet; @@ -74,19 +75,21 @@ public class TsFileSequenceReaderTest { case MetaMarker.CHUNK_HEADER: case MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER: ChunkHeader header = reader.readChunkHeader(marker); - for (int j = 0; j < header.getNumOfPages(); j++) { + int dataSize = header.getDataSize(); + while (dataSize > 0) { PageHeader pageHeader = reader.readPageHeader(header.getDataType(), header.getChunkType() == MetaMarker.CHUNK_HEADER); - reader.readPage(pageHeader, header.getCompressionType()); + ByteBuffer pageData = reader.readPage(pageHeader, header.getCompressionType()); + dataSize -= pageHeader.getSerializedPageSize(); } break; case MetaMarker.CHUNK_GROUP_HEADER: - ChunkGroupHeader footer = reader.readChunkGroupHeader(); + ChunkGroupHeader chunkGroupHeader = reader.readChunkGroupHeader(); long endOffset = reader.position(); Pair<Long, Long> pair = new Pair<>(startOffset, endOffset); - deviceChunkGroupMetadataOffsets.putIfAbsent(footer.getDeviceID(), new ArrayList<>()); + deviceChunkGroupMetadataOffsets.putIfAbsent(chunkGroupHeader.getDeviceID(), new ArrayList<>()); List<Pair<Long, Long>> metadatas = deviceChunkGroupMetadataOffsets - .get(footer.getDeviceID()); + .get(chunkGroupHeader.getDeviceID()); metadatas.add(pair); startOffset = endOffset; break; diff --git a/tsfile/src/test/java/org/apache/iotdb/tsfile/write/TsFileIOWriterTest.java b/tsfile/src/test/java/org/apache/iotdb/tsfile/write/TsFileIOWriterTest.java index 675a77c..a68ba9b 100644 --- a/tsfile/src/test/java/org/apache/iotdb/tsfile/write/TsFileIOWriterTest.java +++ b/tsfile/src/test/java/org/apache/iotdb/tsfile/write/TsFileIOWriterTest.java @@ -87,13 +87,14 @@ public class TsFileIOWriterTest { Assert.assertEquals(TSFileConfig.VERSION_NUMBER, reader.readVersionNumber()); Assert.assertEquals(TSFileConfig.MAGIC_STRING, reader.readTailMagic()); + reader.position(TSFileConfig.MAGIC_STRING.getBytes().length + 1); + // chunk group header Assert.assertEquals(MetaMarker.CHUNK_GROUP_HEADER, reader.readMarker()); - ChunkGroupHeader footer = reader.readChunkGroupHeader(); - Assert.assertEquals(deviceId, footer.getDeviceID()); + ChunkGroupHeader chunkGroupHeader = reader.readChunkGroupHeader(); + Assert.assertEquals(deviceId, chunkGroupHeader.getDeviceID()); // chunk header - reader.position(TSFileConfig.MAGIC_STRING.getBytes().length + 1); Assert.assertEquals(MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER, reader.readMarker()); ChunkHeader header = reader.readChunkHeader(MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER); Assert.assertEquals(TimeSeriesMetadataTest.measurementUID, header.getMeasurementID()); diff --git a/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/PageWriterTest.java b/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/PageWriterTest.java index 3253495..df3a57a 100755 --- a/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/PageWriterTest.java +++ b/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/PageWriterTest.java @@ -46,7 +46,7 @@ public class PageWriterTest { int timeCount = 0; try { writer.write(timeCount++, value); - assertEquals(12, writer.estimateMaxMemSize()); + assertEquals(9, writer.estimateMaxMemSize()); ByteBuffer buffer1 = writer.getUncompressedBytes(); ByteBuffer buffer = ByteBuffer.wrap(buffer1.array()); writer.reset(new MeasurementSchema("s0", TSDataType.INT32, TSEncoding.RLE)); @@ -165,7 +165,7 @@ public class PageWriterTest { int timeCount = 0; try { writer.write(timeCount++, new Binary(value)); - assertEquals(26, writer.estimateMaxMemSize()); + assertEquals(23, writer.estimateMaxMemSize()); ByteBuffer buffer1 = writer.getUncompressedBytes(); ByteBuffer buffer = ByteBuffer.wrap(buffer1.array()); writer.reset(new MeasurementSchema("s0", TSDataType.INT64, TSEncoding.RLE)); diff --git a/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/RestorableTsFileIOWriterTest.java b/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/RestorableTsFileIOWriterTest.java index 68c46dc..0a59614 100644 --- a/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/RestorableTsFileIOWriterTest.java +++ b/tsfile/src/test/java/org/apache/iotdb/tsfile/write/writer/RestorableTsFileIOWriterTest.java @@ -85,8 +85,7 @@ public class RestorableTsFileIOWriterTest { RestorableTsFileIOWriter rWriter = new RestorableTsFileIOWriter(file); writer = new TsFileWriter(rWriter); writer.close(); - assertEquals(TSFileConfig.MAGIC_STRING.getBytes().length + TSFileConfig.VERSION_NUMBER_V2 - .getBytes().length, rWriter.getTruncatedSize()); + assertEquals(TSFileConfig.MAGIC_STRING.getBytes().length + 1, rWriter.getTruncatedSize()); rWriter = new RestorableTsFileIOWriter(file); assertEquals(TsFileCheckStatus.COMPLETE_FILE, rWriter.getTruncatedSize()); @@ -128,6 +127,7 @@ public class RestorableTsFileIOWriterTest { public void testOnlyOneChunkHeader() throws Exception { File file = new File(FILE_NAME); TsFileWriter writer = new TsFileWriter(file); + writer.getIOWriter().startChunkGroup("root.sg1.d1"); writer.getIOWriter() .startFlushChunk(new MeasurementSchema("s1", TSDataType.FLOAT, TSEncoding.PLAIN), CompressionType.SNAPPY, TSDataType.FLOAT, TSEncoding.PLAIN, new FloatStatistics(), 100, @@ -181,6 +181,7 @@ public class RestorableTsFileIOWriterTest { writer.write(new TSRecord(2, "d1").addTuple(new FloatDataPoint("s1", 5)) .addTuple(new FloatDataPoint("s2", 4))); writer.flushAllChunkGroups(); + writer.writeVersion(0); writer.getIOWriter().close(); RestorableTsFileIOWriter rWriter = new RestorableTsFileIOWriter(file); writer = new TsFileWriter(rWriter);
