This is an automated email from the ASF dual-hosted git repository.

jackietien pushed a commit to branch UnseqImprove
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 7ce3c9cd813894a36ed6acc86067080b5b09036c
Author: JackieTien97 <[email protected]>
AuthorDate: Thu Oct 8 21:09:58 2020 +0800

    init
---
 .../iotdb/db/query/reader/series/SeriesReader.java | 178 +++++++++++++++------
 .../iotdb/tsfile/file/metadata/ChunkMetadata.java  |  11 ++
 .../tsfile/file/metadata/TimeseriesMetadata.java   |  11 ++
 3 files changed, 155 insertions(+), 45 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
 
b/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
index e74d4b9..2054f87 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
@@ -95,7 +95,8 @@ public class SeriesReader {
    * page cache
    */
   private VersionPageReader firstPageReader;
-  private PriorityQueue<VersionPageReader> cachedPageReaders;
+  private final List<VersionPageReader> seqPageReaders = new LinkedList<>();
+  private PriorityQueue<VersionPageReader> unseqPageReaders;
 
   /*
    * point cache
@@ -133,7 +134,7 @@ public class SeriesReader {
         timeSeriesMetadata -> 
orderUtils.getOrderTime(timeSeriesMetadata.getStatistics())));
     cachedChunkMetadata = new PriorityQueue<>(orderUtils.comparingLong(
         chunkMetadata -> 
orderUtils.getOrderTime(chunkMetadata.getStatistics())));
-    cachedPageReaders = new PriorityQueue<>(orderUtils.comparingLong(
+    unseqPageReaders = new PriorityQueue<>(orderUtils.comparingLong(
         versionPageReader -> 
orderUtils.getOrderTime(versionPageReader.getStatistics())));
   }
 
@@ -163,7 +164,7 @@ public class SeriesReader {
         timeSeriesMetadata -> 
orderUtils.getOrderTime(timeSeriesMetadata.getStatistics())));
     cachedChunkMetadata = new PriorityQueue<>(orderUtils.comparingLong(
         chunkMetadata -> 
orderUtils.getOrderTime(chunkMetadata.getStatistics())));
-    cachedPageReaders = new PriorityQueue<>(orderUtils.comparingLong(
+    unseqPageReaders = new PriorityQueue<>(orderUtils.comparingLong(
         versionPageReader -> 
orderUtils.getOrderTime(versionPageReader.getStatistics())));
   }
 
@@ -173,12 +174,12 @@ public class SeriesReader {
 
   boolean hasNextFile() throws IOException {
 
-    if (!cachedPageReaders.isEmpty()
+    if (!unseqPageReaders.isEmpty()
         || firstPageReader != null
         || mergeReader.hasNextTimeValuePair()) {
       throw new IOException(
           "all cached pages should be consumed first 
cachedPageReaders.isEmpty() is "
-              + cachedPageReaders.isEmpty()
+              + unseqPageReaders.isEmpty()
               + " firstPageReader != null is "
               + (firstPageReader != null)
               + " mergeReader.hasNextTimeValuePair() = "
@@ -231,12 +232,12 @@ public class SeriesReader {
    * overlapped chunks are consumed
    */
   boolean hasNextChunk() throws IOException {
-    if (!cachedPageReaders.isEmpty()
+    if (!unseqPageReaders.isEmpty()
         || firstPageReader != null
         || mergeReader.hasNextTimeValuePair()) {
       throw new IOException(
           "all cached pages should be consumed first 
cachedPageReaders.isEmpty() is "
-              + cachedPageReaders.isEmpty()
+              + unseqPageReaders.isEmpty()
               + " firstPageReader != null is "
               + (firstPageReader != null)
               + " mergeReader.hasNextTimeValuePair() = "
@@ -300,6 +301,7 @@ public class SeriesReader {
       throws IOException {
     List<ChunkMetadata> chunkMetadataList = FileLoaderUtils
         .loadChunkMetadataList(timeSeriesMetadata);
+    chunkMetadataList.forEach(chunkMetadata -> 
chunkMetadata.setSeq(timeSeriesMetadata.isSeq()));
     // try to calculate the total number of chunk and time-value points in 
chunk
     if 
(IoTDBDescriptor.getInstance().getConfig().isEnablePerformanceTracing()) {
       QueryResourceManager queryResourceManager = 
QueryResourceManager.getInstance();
@@ -381,19 +383,28 @@ public class SeriesReader {
       /*
        * first chunk metadata is already unpacked, consume cached pages
        */
-      if (!cachedPageReaders.isEmpty()) {
-        firstPageReader = cachedPageReaders.poll();
-        long endpointTime = 
orderUtils.getOverlapCheckTime(firstPageReader.getStatistics());
-        unpackAllOverlappedTsFilesToTimeSeriesMetadata(endpointTime);
-        
unpackAllOverlappedTimeSeriesMetadataToCachedChunkMetadata(endpointTime, false);
-        unpackAllOverlappedChunkMetadataToCachedPageReaders(endpointTime, 
false);
+      if (!seqPageReaders.isEmpty() && !unseqPageReaders.isEmpty()) {
+        if (seqPageReaders.get(0).getStatistics().getStartTime() < 
unseqPageReaders.peek()
+            .getStatistics().getStartTime()) {
+          firstPageReader = seqPageReaders.remove(0);
+        } else {
+          firstPageReader = unseqPageReaders.poll();
+        }
+      } else if (!seqPageReaders.isEmpty()) {
+        firstPageReader = seqPageReaders.remove(0);
+      } else {
+        firstPageReader = unseqPageReaders.poll();
       }
+      long endpointTime = 
orderUtils.getOverlapCheckTime(firstPageReader.getStatistics());
+      unpackAllOverlappedTsFilesToTimeSeriesMetadata(endpointTime);
+      unpackAllOverlappedTimeSeriesMetadataToCachedChunkMetadata(endpointTime, 
false);
+      unpackAllOverlappedChunkMetadataToCachedPageReaders(endpointTime, false);
     }
 
-    if (firstPageReader != null
-        && !cachedPageReaders.isEmpty()
-        && orderUtils
-        .isOverlapped(firstPageReader.getStatistics(), 
cachedPageReaders.peek().getStatistics())) {
+    if (firstPageReader != null && ((!unseqPageReaders.isEmpty() && orderUtils
+        .isOverlapped(firstPageReader.getStatistics(), 
unseqPageReaders.peek().getStatistics()))
+        || (!seqPageReaders.isEmpty() && orderUtils
+        .isOverlapped(firstPageReader.getStatistics(), 
seqPageReaders.get(0).getStatistics())))) {
       /*
        * next page is overlapped, read overlapped data and cache it
        */
@@ -407,11 +418,26 @@ public class SeriesReader {
     }
 
     // make sure firstPageReader won't be null while cachedPageReaders has 
more cached page readers
-    while (firstPageReader == null && !cachedPageReaders.isEmpty()) {
-      firstPageReader = cachedPageReaders.poll();
-      if (!cachedPageReaders.isEmpty()
-          && orderUtils.isOverlapped(firstPageReader.getStatistics(),
-          cachedPageReaders.peek().getStatistics())) {
+    while (firstPageReader == null && (!seqPageReaders.isEmpty() || 
!unseqPageReaders.isEmpty())) {
+
+      if (!seqPageReaders.isEmpty() && !unseqPageReaders.isEmpty()) {
+        if (seqPageReaders.get(0).getStatistics().getStartTime() < 
unseqPageReaders.peek()
+            .getStatistics().getStartTime()) {
+          firstPageReader = seqPageReaders.remove(0);
+        } else {
+          firstPageReader = unseqPageReaders.poll();
+        }
+      } else if (!seqPageReaders.isEmpty()) {
+        firstPageReader = seqPageReaders.remove(0);
+      } else {
+        firstPageReader = unseqPageReaders.poll();
+      }
+
+      if ((!seqPageReaders.isEmpty() && orderUtils
+          .isOverlapped(firstPageReader.getStatistics(), 
seqPageReaders.get(0).getStatistics()))
+          || (!unseqPageReaders.isEmpty() && orderUtils
+          .isOverlapped(firstPageReader.getStatistics(),
+              unseqPageReaders.peek().getStatistics()))) {
         /*
          * next page is overlapped, read overlapped data and cache it
          */
@@ -438,17 +464,22 @@ public class SeriesReader {
       unpackOneChunkMetaData(firstChunkMetadata);
       firstChunkMetadata = null;
     }
-    if (init && firstPageReader == null && !cachedPageReaders.isEmpty()) {
-      firstPageReader = cachedPageReaders.poll();
+    if (init && firstPageReader == null && !unseqPageReaders.isEmpty()) {
+      firstPageReader = unseqPageReaders.poll();
     }
   }
 
-  private void unpackOneChunkMetaData(ChunkMetadata chunkMetaData) throws 
IOException {
-    FileLoaderUtils.loadPageReaderList(chunkMetaData, timeFilter)
-        .forEach(
-            pageReader ->
-                cachedPageReaders.add(
-                    new VersionPageReader(chunkMetaData.getVersion(), 
pageReader)));
+  private void unpackOneChunkMetaData(ChunkMetadata chunkMetaData)
+      throws IOException {
+    FileLoaderUtils.loadPageReaderList(chunkMetaData, timeFilter).forEach(
+        pageReader -> {
+          if (chunkMetaData.isSeq()) {
+            seqPageReaders.add(new 
VersionPageReader(chunkMetaData.getVersion(), pageReader, true));
+          } else {
+            unseqPageReaders
+                .add(new VersionPageReader(chunkMetaData.getVersion(), 
pageReader, false));
+          }
+        });
   }
 
   /**
@@ -469,14 +500,16 @@ public class SeriesReader {
     /*
      * has a non-overlapped page in firstPageReader
      */
-    if (mergeReader.hasNextTimeValuePair()) {
+    if (mergeReader.hasNextTimeValuePair()
+        && mergeReader.currentTimeValuePair().getTimestamp() <= 
firstPageReader.getStatistics()
+        .getEndTime()) {
       throw new IOException("overlapped data should be consumed first");
     }
 
     Statistics firstPageStatistics = firstPageReader.getStatistics();
 
-    return !cachedPageReaders.isEmpty()
-        && orderUtils.isOverlapped(firstPageStatistics, 
cachedPageReaders.peek().getStatistics());
+    return !unseqPageReaders.isEmpty()
+        && orderUtils.isOverlapped(firstPageStatistics, 
unseqPageReaders.peek().getStatistics());
   }
 
   Statistics currentPageStatistics() {
@@ -543,7 +576,14 @@ public class SeriesReader {
 
         cachedBatchData = BatchDataFactory.createBatchData(dataType);
         long currentPageEndPointTime = mergeReader.getCurrentReadStopTime();
-
+        if (firstPageReader != null) {
+          currentPageEndPointTime = Math
+              .min(currentPageEndPointTime, 
firstPageReader.getStatistics().getEndTime());
+        }
+        if (!seqPageReaders.isEmpty()) {
+          currentPageEndPointTime = Math
+              .min(currentPageEndPointTime, 
seqPageReaders.get(0).getStatistics().getEndTime());
+        }
         while (mergeReader.hasNextTimeValuePair()) {
 
           /*
@@ -559,7 +599,35 @@ public class SeriesReader {
           unpackAllOverlappedTimeSeriesMetadataToCachedChunkMetadata(
               timeValuePair.getTimestamp(), false);
           
unpackAllOverlappedChunkMetadataToCachedPageReaders(timeValuePair.getTimestamp(),
 false);
-          
unpackAllOverlappedCachedPageReadersToMergeReader(timeValuePair.getTimestamp());
+          
unpackAllOverlappedUnseqPageReadersToMergeReader(timeValuePair.getTimestamp());
+
+          /*
+           * get the latest first point in mergeReader
+           */
+          timeValuePair = mergeReader.currentTimeValuePair();
+
+          if (firstPageReader != null) {
+            if (firstPageReader.getStatistics().getEndTime() < 
timeValuePair.getTimestamp()) {
+              return cachedBatchData.hasCurrent();
+            } else {
+              mergeReader
+                  
.addReader(firstPageReader.getAllSatisfiedPageData(orderUtils.getAscending())
+                          .getBatchDataIterator(), firstPageReader.version,
+                      
orderUtils.getOverlapCheckTime(firstPageReader.getStatistics()));
+              firstPageReader = null;
+            }
+          }
+
+          if (!seqPageReaders.isEmpty()) {
+            if (seqPageReaders.get(0).getStatistics().getEndTime() < 
timeValuePair.getTimestamp()) {
+              return cachedBatchData.hasCurrent();
+            } else {
+              VersionPageReader pageReader = seqPageReaders.remove(0);
+              
mergeReader.addReader(pageReader.getAllSatisfiedPageData(orderUtils.getAscending())
+                      .getBatchDataIterator(), pageReader.version,
+                  orderUtils.getOverlapCheckTime(pageReader.getStatistics()));
+            }
+          }
 
           /*
            * get the latest first point in mergeReader
@@ -591,7 +659,7 @@ public class SeriesReader {
     /*
      * no cached page readers
      */
-    if (firstPageReader == null && cachedPageReaders.isEmpty()) {
+    if (firstPageReader == null && unseqPageReaders.isEmpty() && 
seqPageReaders.isEmpty()) {
       return;
     }
 
@@ -599,7 +667,18 @@ public class SeriesReader {
      * init firstPageReader
      */
     if (firstPageReader == null) {
-      firstPageReader = cachedPageReaders.poll();
+      if (!seqPageReaders.isEmpty() && !unseqPageReaders.isEmpty()) {
+        if (seqPageReaders.get(0).getStatistics().getStartTime() < 
unseqPageReaders.peek()
+            .getStatistics().getStartTime()) {
+          firstPageReader = seqPageReaders.remove(0);
+        } else {
+          firstPageReader = unseqPageReaders.poll();
+        }
+      } else if (!seqPageReaders.isEmpty()) {
+        firstPageReader = seqPageReaders.remove(0);
+      } else {
+        firstPageReader = unseqPageReaders.poll();
+      }
     }
 
     long currentPageEndpointTime;
@@ -610,18 +689,18 @@ public class SeriesReader {
     }
 
     /*
-     * put all currently directly overlapped page reader to merge reader
+     * put all currently directly overlapped unseq page reader to merge reader
      */
-    unpackAllOverlappedCachedPageReadersToMergeReader(currentPageEndpointTime);
+    unpackAllOverlappedUnseqPageReadersToMergeReader(currentPageEndpointTime);
   }
 
-  private void unpackAllOverlappedCachedPageReadersToMergeReader(long 
endpointTime)
+  private void unpackAllOverlappedUnseqPageReadersToMergeReader(long 
endpointTime)
       throws IOException {
-    while (!cachedPageReaders.isEmpty()
-        && orderUtils.isOverlapped(endpointTime, 
cachedPageReaders.peek().data.getStatistics())) {
-      putPageReaderToMergeReader(cachedPageReaders.poll());
+    while (!unseqPageReaders.isEmpty()
+        && orderUtils.isOverlapped(endpointTime, 
unseqPageReaders.peek().data.getStatistics())) {
+      putPageReaderToMergeReader(unseqPageReaders.poll());
     }
-    if (firstPageReader != null &&
+    if (firstPageReader != null && !firstPageReader.isSeq &&
         orderUtils.isOverlapped(endpointTime, 
firstPageReader.getStatistics())) {
       putPageReaderToMergeReader(firstPageReader);
       firstPageReader = null;
@@ -668,6 +747,7 @@ public class SeriesReader {
               orderUtils.getNextSeqFileResource(seqFileResource, true), 
seriesPath, context,
               getAnyFilter(), allSensors);
       if (timeseriesMetadata != null) {
+        timeseriesMetadata.setSeq(true);
         seqTimeSeriesMetadata.add(timeseriesMetadata);
       }
     }
@@ -681,6 +761,7 @@ public class SeriesReader {
               unseqFileResource.remove(0), seriesPath, context, 
getAnyFilter(), allSensors);
       if (timeseriesMetadata != null) {
         timeseriesMetadata.setModified(true);
+        timeseriesMetadata.setSeq(false);
         unSeqTimeSeriesMetadata.add(timeseriesMetadata);
       }
     }
@@ -775,9 +856,12 @@ public class SeriesReader {
     protected long version;
     protected IPageReader data;
 
-    VersionPageReader(long version, IPageReader data) {
+    protected boolean isSeq;
+
+    VersionPageReader(long version, IPageReader data, boolean isSeq) {
       this.version = version;
       this.data = data;
+      this.isSeq = isSeq;
     }
 
     Statistics getStatistics() {
@@ -795,6 +879,10 @@ public class SeriesReader {
     boolean isModified() {
       return data.isModified();
     }
+
+    public boolean isSeq() {
+      return isSeq;
+    }
   }
 
 
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 da68c93..499bf57 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
@@ -73,6 +73,9 @@ public class ChunkMetadata implements Accountable {
   private static final int CHUNK_METADATA_FIXED_RAM_SIZE = 80;
 
 
+  // used for SeriesReader to indicate whether it is a seq/unseq timeseries 
metadata
+  private boolean isSeq = true;
+
   private ChunkMetadata() {
   }
 
@@ -269,4 +272,12 @@ public class ChunkMetadata implements Accountable {
     this.statistics.mergeStatistics(chunkMetadata.getStatistics());
     this.ramSize = calculateRamSize();
   }
+
+  public void setSeq(boolean seq) {
+    isSeq = seq;
+  }
+
+  public boolean isSeq() {
+    return isSeq;
+  }
 }
\ No newline at end of file
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 754d48b..ad8fc3f 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
@@ -46,6 +46,9 @@ public class TimeseriesMetadata implements Accountable {
 
   private long ramSize;
 
+  // used for SeriesReader to indicate whether it is a seq/unseq timeseries 
metadata
+  private boolean isSeq = true;
+
   public TimeseriesMetadata() {
   }
 
@@ -158,4 +161,12 @@ public class TimeseriesMetadata implements Accountable {
   public long getRamSize() {
     return ramSize;
   }
+  
+  public void setSeq(boolean seq) {
+    isSeq = seq;
+  }
+
+  public boolean isSeq() {
+    return isSeq;
+  }
 }

Reply via email to