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;

Reply via email to