This is an automated email from the ASF dual-hosted git repository. jackietien pushed a commit to branch AlignedBug in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 83b4188ad78ac2667203c888bc3733b58d15ecb5 Author: JackieTien97 <[email protected]> AuthorDate: Fri May 27 22:00:26 2022 +0800 fix aligned page reader bug --- .../db/mpp/execution/operator/source/SeriesScanUtil.java | 3 +++ .../db/query/reader/chunk/MemAlignedPageReader.java | 16 +++++++++------- .../iotdb/db/query/reader/chunk/MemPageReader.java | 4 ++++ .../org/apache/iotdb/tsfile/read/reader/IPageReader.java | 4 ++++ .../iotdb/tsfile/read/reader/page/AlignedPageReader.java | 13 +++++++------ .../apache/iotdb/tsfile/read/reader/page/PageReader.java | 3 +++ 6 files changed, 30 insertions(+), 13 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/source/SeriesScanUtil.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/source/SeriesScanUtil.java index 5bde17b84b..279f6be932 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/source/SeriesScanUtil.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/source/SeriesScanUtil.java @@ -514,6 +514,9 @@ public class SeriesScanUtil { List<IPageReader> pageReaderList = FileLoaderUtils.loadPageReaderList(chunkMetaData, timeFilter); + // init TsBlockBuilder for each page reader + pageReaderList.forEach(p -> p.initTsBlockBuilder(getTsDataTypeList())); + if (chunkMetaData.isSeq()) { if (orderUtils.getAscending()) { for (IPageReader iPageReader : pageReaderList) { diff --git a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java index 30162e7c82..a29d20bf12 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java +++ b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java @@ -19,7 +19,6 @@ package org.apache.iotdb.db.query.reader.chunk; import org.apache.iotdb.tsfile.file.metadata.AlignedChunkMetadata; -import org.apache.iotdb.tsfile.file.metadata.IChunkMetadata; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics; import org.apache.iotdb.tsfile.read.common.BatchData; @@ -35,7 +34,7 @@ import org.apache.iotdb.tsfile.read.reader.IPageReader; import org.apache.iotdb.tsfile.utils.TsPrimitiveType; import java.io.IOException; -import java.util.stream.Collectors; +import java.util.List; public class MemAlignedPageReader implements IPageReader, IAlignedPageReader { @@ -43,6 +42,8 @@ public class MemAlignedPageReader implements IPageReader, IAlignedPageReader { private final AlignedChunkMetadata chunkMetadata; private Filter valueFilter; + private TsBlockBuilder builder; + public MemAlignedPageReader(TsBlock tsBlock, AlignedChunkMetadata chunkMetadata, Filter filter) { this.tsBlock = tsBlock; this.chunkMetadata = chunkMetadata; @@ -87,11 +88,7 @@ public class MemAlignedPageReader implements IPageReader, IAlignedPageReader { @Override public TsBlock getAllSatisfiedData() { - TsBlockBuilder builder = - new TsBlockBuilder( - chunkMetadata.getValueChunkMetadataList().stream() - .map(IChunkMetadata::getDataType) - .collect(Collectors.toList())); + builder.reset(); boolean[] satisfyInfo = new boolean[tsBlock.getPositionCount()]; @@ -158,4 +155,9 @@ public class MemAlignedPageReader implements IPageReader, IAlignedPageReader { public boolean isModified() { return false; } + + @Override + public void initTsBlockBuilder(List<TSDataType> dataTypes) { + builder = new TsBlockBuilder(dataTypes); + } } diff --git a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemPageReader.java b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemPageReader.java index 9d1d9822e1..0baf315eff 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemPageReader.java +++ b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemPageReader.java @@ -35,6 +35,7 @@ import org.apache.iotdb.tsfile.utils.Binary; import java.io.IOException; import java.util.Collections; +import java.util.List; public class MemPageReader implements IPageReader { @@ -184,4 +185,7 @@ public class MemPageReader implements IPageReader { public boolean isModified() { return false; } + + @Override + public void initTsBlockBuilder(List<TSDataType> dataTypes) {} } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IPageReader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IPageReader.java index 3d7db3ec21..a68f4590b1 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IPageReader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IPageReader.java @@ -18,12 +18,14 @@ */ package org.apache.iotdb.tsfile.read.reader; +import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics; import org.apache.iotdb.tsfile.read.common.BatchData; import org.apache.iotdb.tsfile.read.common.block.TsBlock; import org.apache.iotdb.tsfile.read.filter.basic.Filter; import java.io.IOException; +import java.util.List; public interface IPageReader { @@ -40,4 +42,6 @@ public interface IPageReader { void setFilter(Filter filter); boolean isModified(); + + void initTsBlockBuilder(List<TSDataType> dataTypes); } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java index 5dc9a466a6..89893a781c 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java @@ -37,7 +37,6 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; -import java.util.stream.Collectors; public class AlignedPageReader implements IPageReader, IAlignedPageReader { @@ -46,6 +45,7 @@ public class AlignedPageReader implements IPageReader, IAlignedPageReader { private final int valueCount; private Filter filter; private boolean isModified; + private TsBlockBuilder builder; public AlignedPageReader( PageHeader timePageHeader, @@ -108,11 +108,7 @@ public class AlignedPageReader implements IPageReader, IAlignedPageReader { @Override public TsBlock getAllSatisfiedData() throws IOException { // TODO change from the row-based style to column-based style - TsBlockBuilder builder = - new TsBlockBuilder( - valuePageReaderList.stream() - .map(ValuePageReader::getDataType) - .collect(Collectors.toList())); + builder.reset(); int timeIndex = -1; while (timePageReader.hasNextTime()) { long timestamp = timePageReader.nextTime(); @@ -185,4 +181,9 @@ public class AlignedPageReader implements IPageReader, IAlignedPageReader { public boolean isModified() { return isModified; } + + @Override + public void initTsBlockBuilder(List<TSDataType> dataTypes) { + builder = new TsBlockBuilder(dataTypes); + } } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/PageReader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/PageReader.java index b54278451a..e1fce8ff65 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/PageReader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/PageReader.java @@ -265,6 +265,9 @@ public class PageReader implements IPageReader { return pageHeader.isModified(); } + @Override + public void initTsBlockBuilder(List<TSDataType> dataTypes) {} + protected boolean isDeleted(long timestamp) { while (deleteIntervalList != null && deleteCursor < deleteIntervalList.size()) { if (deleteIntervalList.get(deleteCursor).contains(timestamp)) {
