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 {}
+}