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 d5c34e8d2e378513e765779968c3b15964e9a32f Author: JackieTien97 <[email protected]> AuthorDate: Tue Dec 1 21:05:43 2020 +0800 fuck bug day --- .../main/java/org/apache/iotdb/JDBCExample.java | 26 ++++++---- .../apache/iotdb/tsfile/TsFileSequenceRead.java | 3 +- .../engine/merge/MaxFileMergeFileSelectorTest.java | 4 +- .../merge/MaxSeriesMergeFileSelectorTest.java | 8 ++-- .../tsfile/encoding/encoder/PlainEncoder.java | 6 ++- .../iotdb/tsfile/file/header/ChunkHeader.java | 18 +++++-- .../metadata/statistics/BooleanStatistics.java | 5 +- .../org/apache/iotdb/tsfile/read/common/Chunk.java | 56 +++++++++++++++++++--- .../tsfile/read/reader/chunk/ChunkReader.java | 13 ++++- .../iotdb/tsfile/write/chunk/ChunkWriterImpl.java | 2 + 10 files changed, 107 insertions(+), 34 deletions(-) diff --git a/example/jdbc/src/main/java/org/apache/iotdb/JDBCExample.java b/example/jdbc/src/main/java/org/apache/iotdb/JDBCExample.java index 00f1084..e8c05ce 100644 --- a/example/jdbc/src/main/java/org/apache/iotdb/JDBCExample.java +++ b/example/jdbc/src/main/java/org/apache/iotdb/JDBCExample.java @@ -28,21 +28,28 @@ import java.sql.SQLException; import java.sql.Statement; public class JDBCExample { + public static void main(String[] args) throws ClassNotFoundException, SQLException { Class.forName("org.apache.iotdb.jdbc.IoTDBDriver"); - try (Connection connection = DriverManager.getConnection("jdbc:iotdb://127.0.0.1:6667/", "root", "root"); - Statement statement = connection.createStatement()) { + try (Connection connection = DriverManager + .getConnection("jdbc:iotdb://127.0.0.1:6667/", "root", "root"); + Statement statement = connection.createStatement()) { try { statement.execute("SET STORAGE GROUP TO root.sg1"); - statement.execute("CREATE TIMESERIES root.sg1.d1.s1 WITH DATATYPE=INT64, ENCODING=RLE, COMPRESSOR=SNAPPY"); - statement.execute("CREATE TIMESERIES root.sg1.d1.s2 WITH DATATYPE=INT64, ENCODING=RLE, COMPRESSOR=SNAPPY"); - statement.execute("CREATE TIMESERIES root.sg1.d1.s3 WITH DATATYPE=INT64, ENCODING=RLE, COMPRESSOR=SNAPPY"); + statement.execute( + "CREATE TIMESERIES root.sg1.d1.s1 WITH DATATYPE=INT64, ENCODING=RLE, COMPRESSOR=SNAPPY"); + statement.execute( + "CREATE TIMESERIES root.sg1.d1.s2 WITH DATATYPE=INT64, ENCODING=RLE, COMPRESSOR=SNAPPY"); + statement.execute( + "CREATE TIMESERIES root.sg1.d1.s3 WITH DATATYPE=INT64, ENCODING=RLE, COMPRESSOR=SNAPPY"); } catch (IoTDBSQLException e) { System.out.println(e.getMessage()); } for (int i = 0; i <= 100; i++) { - statement.addBatch("insert into root.sg1.d1(timestamp, s1, s2, s3) values("+ i + "," + 1 + "," + 1 + "," + 1 + ")"); + statement.addBatch( + "insert into root.sg1.d1(timestamp, s1, s2, s3) values(" + i + "," + 1 + "," + 1 + "," + + 1 + ")"); } statement.executeBatch(); statement.clearBatch(); @@ -51,10 +58,11 @@ public class JDBCExample { outputResult(resultSet); resultSet = statement.executeQuery("select count(*) from root"); outputResult(resultSet); - resultSet = statement.executeQuery("select count(*) from root where time >= 1 and time <= 100 group by ([0, 100), 20ms, 20ms)"); + resultSet = statement.executeQuery( + "select count(*) from root where time >= 1 and time <= 100 group by ([0, 100), 20ms, 20ms)"); outputResult(resultSet); - } catch (IoTDBSQLException e){ - System.out.println(e.getMessage()); + } catch (IoTDBSQLException e) { + System.out.println(e.getMessage()); } } 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 93ed9f0..a3ccda0 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,7 +41,7 @@ public class TsFileSequenceRead { @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity warning public static void main(String[] args) throws IOException { - String filename = "/Users/jackietien/Desktop/1-1-1-after.tsfile"; + String filename = "test.tsfile"; if (args.length >= 1) { filename = args[0]; } @@ -65,6 +65,7 @@ public class TsFileSequenceRead { case MetaMarker.CHUNK_HEADER: case MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER: System.out.println("\t[Chunk]"); + System.out.println("\tchunk type: " + marker); System.out.println("\tposition: " + reader.position()); ChunkHeader header = reader.readChunkHeader(marker); System.out.println("\tMeasurement: " + header.getMeasurementID()); diff --git a/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java b/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java index 1e50590..7036f75 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java @@ -78,8 +78,8 @@ public class MaxFileMergeFileSelectorTest extends MergeTest { List[] result = mergeFileSelector.select(); List<TsFileResource> seqSelected = result[0]; List<TsFileResource> unseqSelected = result[1]; - assertEquals(seqResources.subList(0, 3), seqSelected); - assertEquals(unseqResources.subList(0, 3), unseqSelected); + assertEquals(seqResources.subList(0, 4), seqSelected); + assertEquals(unseqResources.subList(0, 4), unseqSelected); resource.clear(); } } diff --git a/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxSeriesMergeFileSelectorTest.java b/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxSeriesMergeFileSelectorTest.java index 2a876d1..2e89ee5 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxSeriesMergeFileSelectorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxSeriesMergeFileSelectorTest.java @@ -85,8 +85,8 @@ public class MaxSeriesMergeFileSelectorTest extends MergeTest { List[] result = mergeFileSelector.select(); List<TsFileResource> seqSelected = result[0]; List<TsFileResource> unseqSelected = result[1]; - assertEquals(seqResources.subList(0, 3), seqSelected); - assertEquals(unseqResources.subList(0, 3), unseqSelected); + assertEquals(seqResources.subList(0, 4), seqSelected); + assertEquals(unseqResources.subList(0, 4), unseqSelected); assertEquals(MaxSeriesMergeFileSelector.MAX_SERIES_NUM, mergeFileSelector.getConcurrentMergeNum()); resource.clear(); @@ -100,8 +100,8 @@ public class MaxSeriesMergeFileSelectorTest extends MergeTest { List[] result = mergeFileSelector.select(); List<TsFileResource> seqSelected = result[0]; List<TsFileResource> unseqSelected = result[1]; - assertEquals(seqResources.subList(0, 1), seqSelected); - assertEquals(unseqResources.subList(0, 1), unseqSelected); + assertEquals(seqResources.subList(0, 2), seqSelected); + assertEquals(unseqResources.subList(0, 2), unseqSelected); resource.clear(); } } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/encoder/PlainEncoder.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/encoder/PlainEncoder.java index 50bff74..5fe4ccd 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/encoder/PlainEncoder.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/encoder/PlainEncoder.java @@ -72,7 +72,11 @@ public class PlainEncoder extends Encoder { @Override public void encode(float value, ByteArrayOutputStream out) { - encode(Float.floatToIntBits(value), out); + int floatInt = Float.floatToIntBits(value); + out.write((floatInt >> 24) & 0xFF); + out.write((floatInt >> 16) & 0xFF); + out.write((floatInt >> 8) & 0xFF); + out.write(floatInt & 0xFF); } @Override 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 04942d8..082a4c7 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 @@ -35,6 +35,11 @@ import java.nio.ByteBuffer; public class ChunkHeader { + /** + * 1 means this chunk has more than one page, so each page has its own page statistic 4 means this + * chunk has only one page, and this page has no page statistic + */ + private byte chunkType; private String measurementID; private int dataSize; private TSDataType dataType; @@ -42,11 +47,6 @@ public class ChunkHeader { private TSEncoding encodingType; // the following fields do not need to be serialized. - /** - * 1 means this chunk has more than one page, so each page has its own page statistic 4 means this - * chunk has only one page, and this page has no page statistic - */ - private byte chunkType; private int numOfPages; private int serializedSize; @@ -226,10 +226,18 @@ public class ChunkHeader { this.numOfPages += chunkHeader.getNumOfPages(); } + public void setDataSize(int dataSize) { + this.dataSize = dataSize; + } + public byte getChunkType() { return chunkType; } + public void setChunkType(byte chunkType) { + this.chunkType = chunkType; + } + public void increasePageNums(int i) { numOfPages += i; } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/BooleanStatistics.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/BooleanStatistics.java index ddbdf34..201052f 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/BooleanStatistics.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/BooleanStatistics.java @@ -220,10 +220,9 @@ public class BooleanStatistics extends Statistics<Boolean> { @Override public String toString() { - return "BooleanStatistics{" + - "firstValue=" + firstValue + + return super.toString() + " [firstValue=" + firstValue + ", lastValue=" + lastValue + ", sumValue=" + sumValue + - '}'; + ']'; } } 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 a7bbadb..58d15f8 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 @@ -18,12 +18,15 @@ */ package org.apache.iotdb.tsfile.read.common; +import java.io.IOException; import java.nio.ByteBuffer; - import java.util.List; import org.apache.iotdb.tsfile.common.cache.Accountable; +import org.apache.iotdb.tsfile.file.MetaMarker; import org.apache.iotdb.tsfile.file.header.ChunkHeader; import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics; +import org.apache.iotdb.tsfile.utils.PublicBAOS; +import org.apache.iotdb.tsfile.utils.ReadWriteForEncodingUtils; /** * used in query. @@ -63,12 +66,51 @@ public class Chunk implements Accountable { this.deleteIntervalList = list; } - public void mergeChunk(Chunk chunk) { - chunkHeader.mergeChunkHeader(chunk.chunkHeader); - ByteBuffer newChunkData = ByteBuffer - .allocate(chunkData.array().length + chunk.chunkData.array().length); - newChunkData.put(chunkData.array()); - newChunkData.put(chunk.chunkData.array()); + public void mergeChunk(Chunk chunk) throws IOException { + int dataSize = 0; + int offset1 = -1; + if (chunk.chunkHeader.getChunkType() == MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER) { + ReadWriteForEncodingUtils.readUnsignedVarInt(chunk.chunkData); + ReadWriteForEncodingUtils.readUnsignedVarInt(chunk.chunkData); + offset1 = chunk.chunkData.position(); + chunk.chunkData.flip(); + dataSize += (chunk.chunkData.array().length + chunk.chunkStatistic.getSerializedSize()); + } else { + dataSize += chunk.chunkData.array().length; + } + int offset2 = -1; + if (chunkHeader.getChunkType() == MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER) { + chunkHeader.setChunkType(MetaMarker.CHUNK_HEADER); + ReadWriteForEncodingUtils.readUnsignedVarInt(chunkData); + ReadWriteForEncodingUtils.readUnsignedVarInt(chunkData); + offset2 = chunkData.position(); + chunkData.flip(); + dataSize += (chunkData.array().length + chunkStatistic.getSerializedSize()); + } else { + dataSize += chunkData.array().length; + } + chunkHeader.setDataSize(dataSize); + ByteBuffer newChunkData = ByteBuffer.allocate(dataSize); + if (offset2 == -1) { + newChunkData.put(chunkData.array()); + } else { + byte[] b = chunkData.array(); + newChunkData.put(b, 0, offset2); + PublicBAOS a = new PublicBAOS(); + chunkStatistic.serialize(a); + newChunkData.put(a.getBuf(), 0, a.size()); + newChunkData.put(b, offset2, b.length - offset2); + } + if (offset1 == -1) { + newChunkData.put(chunk.chunkData.array()); + } else { + byte[] b = chunk.chunkData.array(); + newChunkData.put(b, 0, offset1); + PublicBAOS a = new PublicBAOS(); + chunk.chunkStatistic.serialize(a); + newChunkData.put(a.getBuf(), 0, a.size()); + newChunkData.put(b, offset1, b.length - offset1); + } chunkData = newChunkData; } 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 2250307..27a2518 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 @@ -163,8 +163,17 @@ public class ChunkReader implements IChunkReader { Decoder valueDecoder = Decoder .getDecoderByType(chunkHeader.getEncodingType(), chunkHeader.getDataType()); byte[] uncompressedPageData = new byte[pageHeader.getUncompressedSize()]; - unCompressor.uncompress(compressedPageBody, 0, compressedPageBodyLength, - uncompressedPageData, 0); + try { + unCompressor.uncompress(compressedPageBody, 0, compressedPageBodyLength, + uncompressedPageData, 0); + } catch (Exception e) { + System.out.println("error: "); + System.out.println("uncompress size: " + pageHeader.getUncompressedSize()); + System.out.println("compressed size: " + pageHeader.getCompressedSize()); + System.out.println("page header: " + pageHeader); + e.printStackTrace(); + } + ByteBuffer pageData = ByteBuffer.wrap(uncompressedPageData); PageReader reader = new PageReader(pageHeader, pageData, chunkHeader.getDataType(), valueDecoder, timeDecoder, filter); 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 44c2dae..95c4c9d 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 @@ -237,6 +237,8 @@ public class ChunkWriterImpl implements IChunkWriter { // reinit this chunk writer pageBuffer.reset(); + numOfPages = 0; + firstPageStatistics = null; this.statistics = Statistics.getStatsByType(measurementSchema.getType()); }
