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

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


The following commit(s) were added to refs/heads/lazy_page_reader_in_compaction 
by this push:
     new 05e605e4555 add lazy point reader for compaction
05e605e4555 is described below

commit 05e605e4555cec16402393ce01e8b94c3aa86161
Author: Jinrui.Zhang <[email protected]>
AuthorDate: Thu Aug 24 17:32:21 2023 +0800

    add lazy point reader for compaction
---
 .../utils/executor/fast/element/PageElement.java   |  7 +-
 .../utils/executor/fast/element/PointElement.java  |  2 +-
 .../read/reader/chunk/AlignedChunkReader.java      |  6 +-
 .../tsfile/read/reader/page/AlignedPageReader.java |  5 ++
 .../page/LazyLoadAlignedPagePointReader.java       | 89 ++++++++++++++++++++++
 5 files changed, 103 insertions(+), 6 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PageElement.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PageElement.java
index 0977c9d6a44..46e37df479a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PageElement.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PageElement.java
@@ -22,6 +22,7 @@ package 
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.ex
 import org.apache.iotdb.tsfile.file.header.PageHeader;
 import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 import org.apache.iotdb.tsfile.read.reader.IChunkReader;
+import org.apache.iotdb.tsfile.read.reader.IPointReader;
 import org.apache.iotdb.tsfile.read.reader.chunk.AlignedChunkReader;
 import org.apache.iotdb.tsfile.read.reader.chunk.ChunkReader;
 
@@ -38,6 +39,8 @@ public class PageElement {
 
   public TsBlock batchData;
 
+  public IPointReader pointReader;
+
   // compressed page data
   public ByteBuffer pageData;
 
@@ -99,9 +102,9 @@ public class PageElement {
 
   public void deserializePage() throws IOException {
     if (iChunkReader instanceof AlignedChunkReader) {
-      this.batchData =
+      this.pointReader =
           ((AlignedChunkReader) iChunkReader)
-              .readPageData(pageHeader, valuePageHeaders, pageData, 
valuePageDatas);
+              .getPagePointReader(pageHeader, valuePageHeaders, pageData, 
valuePageDatas);
     } else {
       this.batchData = ((ChunkReader) iChunkReader).readPageData(pageHeader, 
pageData);
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PointElement.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PointElement.java
index 94f9bb0d031..8ad82c66f5e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PointElement.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/fast/element/PointElement.java
@@ -39,7 +39,7 @@ public class PointElement {
     if (pageElement.iChunkReader instanceof ChunkReader) {
       this.pointReader = 
pageElement.batchData.getTsBlockSingleColumnIterator();
     } else {
-      this.pointReader = pageElement.batchData.getTsBlockAlignedRowIterator();
+      this.pointReader = pageElement.pointReader;
     }
     this.timeValuePair = pointReader.nextTimeValuePair();
     this.timestamp = timeValuePair.getTimestamp();
diff --git 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/AlignedChunkReader.java
 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/AlignedChunkReader.java
index b7c79f7a1a4..bbbb5fc8864 100644
--- 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/AlignedChunkReader.java
+++ 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/chunk/AlignedChunkReader.java
@@ -31,10 +31,10 @@ import 
org.apache.iotdb.tsfile.file.metadata.statistics.Statistics;
 import org.apache.iotdb.tsfile.read.common.BatchData;
 import org.apache.iotdb.tsfile.read.common.Chunk;
 import org.apache.iotdb.tsfile.read.common.TimeRange;
-import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 import org.apache.iotdb.tsfile.read.reader.IChunkReader;
 import org.apache.iotdb.tsfile.read.reader.IPageReader;
+import org.apache.iotdb.tsfile.read.reader.IPointReader;
 import org.apache.iotdb.tsfile.read.reader.page.AlignedPageReader;
 
 import java.io.IOException;
@@ -277,7 +277,7 @@ public class AlignedChunkReader implements IChunkReader {
   }
 
   /** Read data from compressed page data. Uncompress the page and decode it 
to tsblock data. */
-  public TsBlock readPageData(
+  public IPointReader getPagePointReader(
       PageHeader timePageHeader,
       List<PageHeader> valuePageHeaders,
       ByteBuffer compressedTimePageData,
@@ -323,7 +323,7 @@ public class AlignedChunkReader implements IChunkReader {
             false);
     alignedPageReader.initTsBlockBuilder(valueTypes);
     alignedPageReader.setDeleteIntervalList(valueDeleteIntervalList);
-    return alignedPageReader.getAllSatisfiedData();
+    return alignedPageReader.getLazyPointReader();
   }
 
   private ByteBuffer uncompressPageData(
diff --git 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
index 608a3e8fcb0..b8d9086fd9e 100644
--- 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
+++ 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
@@ -32,6 +32,7 @@ import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 import org.apache.iotdb.tsfile.read.filter.operator.AndFilter;
 import org.apache.iotdb.tsfile.read.reader.IAlignedPageReader;
 import org.apache.iotdb.tsfile.read.reader.IPageReader;
+import org.apache.iotdb.tsfile.read.reader.IPointReader;
 import org.apache.iotdb.tsfile.read.reader.series.PaginationController;
 import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
@@ -160,6 +161,10 @@ public class AlignedPageReader implements IPageReader, 
IAlignedPageReader {
     }
   }
 
+  public IPointReader getLazyPointReader() {
+    return new LazyLoadAlignedPagePointReader(timePageReader, 
valuePageReaderList);
+  }
+
   @Override
   public TsBlock getAllSatisfiedData() throws IOException {
     builder.reset();
diff --git 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/LazyLoadAlignedPagePointReader.java
 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/LazyLoadAlignedPagePointReader.java
new file mode 100644
index 00000000000..da3b6fcbe6b
--- /dev/null
+++ 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/LazyLoadAlignedPagePointReader.java
@@ -0,0 +1,89 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.tsfile.read.reader.page;
+
+import org.apache.iotdb.tsfile.read.TimeValuePair;
+import org.apache.iotdb.tsfile.read.reader.IPointReader;
+import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
+
+import java.io.IOException;
+import java.util.List;
+
+public class LazyLoadAlignedPagePointReader implements IPointReader {
+
+  private TimePageReader timeReader;
+  private List<ValuePageReader> valueReaders;
+
+  private boolean hasNextRow = false;
+
+  private int timeIndex;
+  private long currentTime;
+  private TsPrimitiveType currentRow;
+
+  public LazyLoadAlignedPagePointReader(
+      TimePageReader timeReader, List<ValuePageReader> valueReaders) {
+    this.timeIndex = -1;
+    this.timeReader = timeReader;
+    this.valueReaders = valueReaders;
+  }
+
+  private void prepareNextRow() throws IOException {
+    while (true) {
+      if (!timeReader.hasNextTime()) {
+        hasNextRow = false;
+        return;
+      }
+      currentTime = timeReader.nextTime();
+      timeIndex++;
+      boolean someValueNotNull = false;
+      TsPrimitiveType[] valuesInThisRow = new 
TsPrimitiveType[valueReaders.size()];
+      for (int i = 0; i < valueReaders.size(); i++) {
+        TsPrimitiveType value = valueReaders.get(i).nextValue(currentTime, 
timeIndex);
+        someValueNotNull = someValueNotNull || (value != null);
+        valuesInThisRow[i] = value;
+      }
+      if (someValueNotNull) {
+        currentRow = new TsPrimitiveType.TsVector(valuesInThisRow);
+        hasNextRow = true;
+        break;
+      }
+    }
+  }
+
+  @Override
+  public boolean hasNextTimeValuePair() throws IOException {
+    return hasNextRow;
+  }
+
+  @Override
+  public TimeValuePair nextTimeValuePair() throws IOException {
+    TimeValuePair ret = currentTimeValuePair();
+    prepareNextRow();
+    return ret;
+  }
+
+  @Override
+  public TimeValuePair currentTimeValuePair() throws IOException {
+    return new TimeValuePair(currentTime, currentRow);
+  }
+
+  @Override
+  public void close() throws IOException {}
+}

Reply via email to