This is an automated email from the ASF dual-hosted git repository. shuwenwei pushed a commit to branch fixBug0902 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 05e76debf47f6386c11ede4fdbfb706ab44fde5d Author: shuwenwei <[email protected]> AuthorDate: Tue Sep 2 16:21:43 2025 +0800 Concurrently querying and writing to the MemTable may cause the query results out of order --- .../memtable/AlignedReadOnlyMemChunk.java | 32 +++----- .../dataregion/memtable/ReadOnlyMemChunk.java | 20 ++--- .../db/utils/datastructure/AlignedTVList.java | 5 +- .../datastructure/MemPointIteratorFactory.java | 91 +++++++++++++++++----- .../MergeSortMultiAlignedTVListIterator.java | 2 + .../MergeSortMultiTVListIterator.java | 2 + .../datastructure/MultiAlignedTVListIterator.java | 6 +- .../utils/datastructure/MultiTVListIterator.java | 22 +++++- .../OrderedMultiAlignedTVListIterator.java | 2 + .../datastructure/OrderedMultiTVListIterator.java | 2 + .../iotdb/db/utils/datastructure/TVList.java | 3 + .../memtable/AlignedTVListIteratorTest.java | 8 +- .../memtable/NonAlignedTVListIteratorTest.java | 8 +- 13 files changed, 133 insertions(+), 70 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java index 53a8f8cda2b..e69a5b8b8c3 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java @@ -51,7 +51,6 @@ import java.io.Serializable; import java.util.ArrayList; import java.util.List; import java.util.Map; -import java.util.stream.Collectors; public class AlignedReadOnlyMemChunk extends ReadOnlyMemChunk { private final String timeChunkName; @@ -377,23 +376,7 @@ public class AlignedReadOnlyMemChunk extends ReadOnlyMemChunk { } private void writeValidValuesIntoTsBlock(TsBlockBuilder builder) throws IOException { - List<AlignedTVList> alignedTvLists = - alignedTvListQueryMap.keySet().stream() - .map(x -> (AlignedTVList) x) - .collect(Collectors.toList()); - MemPointIterator timeValuePairIterator = - MemPointIteratorFactory.create( - dataTypes, - columnIndexList, - alignedTvLists, - Ordering.ASC, - null, - timeColumnDeletion, - valueColumnsDeletionList, - floatPrecision, - encodingList, - context.isIgnoreAllNullRows(), - MAX_NUMBER_OF_POINTS_IN_PAGE); + MemPointIterator timeValuePairIterator = createMemPointIterator(Ordering.ASC, null); while (timeValuePairIterator.hasNextTimeValuePair()) { TimeValuePair tvPair = timeValuePairIterator.nextTimeValuePair(); @@ -475,14 +458,17 @@ public class AlignedReadOnlyMemChunk extends ReadOnlyMemChunk { @Override public MemPointIterator createMemPointIterator(Ordering scanOrder, Filter globalTimeFilter) { - List<AlignedTVList> alignedTvLists = - alignedTvListQueryMap.keySet().stream() - .map(x -> (AlignedTVList) x) - .collect(Collectors.toList()); + List<AlignedTVList> tvLists = new ArrayList<>(alignedTvListQueryMap.size()); + List<Integer> tvListRowCounts = new ArrayList<>(alignedTvListQueryMap.size()); + for (Map.Entry<TVList, Integer> entry : alignedTvListQueryMap.entrySet()) { + tvLists.add((AlignedTVList) entry.getKey()); + tvListRowCounts.add(entry.getValue()); + } return MemPointIteratorFactory.create( dataTypes, columnIndexList, - alignedTvLists, + tvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java index 45ef1c75d00..f2fe63c02c6 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java @@ -292,17 +292,7 @@ public class ReadOnlyMemChunk { // read all data in memory chunk and write to tsblock private void writeValidValuesIntoTsBlock(TsBlockBuilder builder) throws IOException { - List<TVList> tvLists = new ArrayList<>(tvListQueryMap.keySet()); - MemPointIterator timeValuePairIterator = - MemPointIteratorFactory.create( - getDataType(), - tvLists, - Ordering.ASC, - null, - deletionList, - floatPrecision, - encoding, - MAX_NUMBER_OF_POINTS_IN_PAGE); + MemPointIterator timeValuePairIterator = createMemPointIterator(Ordering.ASC, null); while (timeValuePairIterator.hasNextTimeValuePair()) { TimeValuePair tvPair = timeValuePairIterator.nextTimeValuePair(); @@ -379,10 +369,16 @@ public class ReadOnlyMemChunk { } public MemPointIterator createMemPointIterator(Ordering scanOrder, Filter globalTimeFilter) { - List<TVList> tvLists = new ArrayList<>(tvListQueryMap.keySet()); + List<TVList> tvLists = new ArrayList<>(tvListQueryMap.size()); + List<Integer> tvListRowCounts = new ArrayList<>(tvListQueryMap.size()); + for (Map.Entry<TVList, Integer> entry : tvListQueryMap.entrySet()) { + tvLists.add(entry.getKey()); + tvListRowCounts.add(entry.getValue()); + } return MemPointIteratorFactory.create( dataType, tvLists, + tvListRowCounts, scanOrder, globalTimeFilter, deletionList, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java index b7474d1583b..81615904b95 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java @@ -1601,6 +1601,7 @@ public abstract class AlignedTVList extends TVList { public AlignedTVListIterator iterator( Ordering scanOrder, + int rowCount, Filter globalTimeFilter, List<TSDataType> dataTypeList, List<Integer> columnIndexList, @@ -1612,6 +1613,7 @@ public abstract class AlignedTVList extends TVList { int maxNumberOfPointsInPage) { return new AlignedTVListIterator( scanOrder, + rowCount, globalTimeFilter, dataTypeList, columnIndexList, @@ -1643,6 +1645,7 @@ public abstract class AlignedTVList extends TVList { public AlignedTVListIterator( Ordering scanOrder, + int rowCount, Filter globalTimeFilter, List<TSDataType> dataTypeList, List<Integer> columnIndexList, @@ -1652,7 +1655,7 @@ public abstract class AlignedTVList extends TVList { List<TSEncoding> encodingList, boolean ignoreAllNullRows, int maxNumberOfPointsInPage) { - super(scanOrder, globalTimeFilter, null, null, null, maxNumberOfPointsInPage); + super(scanOrder, rowCount, globalTimeFilter, null, null, null, maxNumberOfPointsInPage); this.dataTypeList = dataTypeList; this.columnIndexList = (columnIndexList == null) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MemPointIteratorFactory.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MemPointIteratorFactory.java index fa78fac300b..20b82c338c8 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MemPointIteratorFactory.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MemPointIteratorFactory.java @@ -35,18 +35,29 @@ public class MemPointIteratorFactory { // TVListIterator private static MemPointIterator single(List<TVList> tvLists, int maxNumberOfPointsInPage) { - return tvLists.get(0).iterator(Ordering.ASC, null, null, null, null, maxNumberOfPointsInPage); + return tvLists + .get(0) + .iterator( + Ordering.ASC, tvLists.get(0).rowCount, null, null, null, null, maxNumberOfPointsInPage); } private static MemPointIterator single( List<TVList> tvLists, List<TimeRange> deletionList, int maxNumberOfPointsInPage) { return tvLists .get(0) - .iterator(Ordering.ASC, null, deletionList, null, null, maxNumberOfPointsInPage); + .iterator( + Ordering.ASC, + tvLists.get(0).rowCount, + null, + deletionList, + null, + null, + maxNumberOfPointsInPage); } private static MemPointIterator single( List<TVList> tvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> deletionList, @@ -57,6 +68,7 @@ public class MemPointIteratorFactory { .get(0) .iterator( scanOrder, + tvListRowCounts.get(0), globalTimeFilter, deletionList, floatPrecision, @@ -68,7 +80,7 @@ public class MemPointIteratorFactory { private static MemPointIterator mergeSort( TSDataType tsDataType, List<TVList> tvLists, int maxNumberOfPointsInPage) { return new MergeSortMultiTVListIterator( - Ordering.ASC, null, tsDataType, tvLists, null, null, null, maxNumberOfPointsInPage); + Ordering.ASC, null, tsDataType, tvLists, null, null, null, null, maxNumberOfPointsInPage); } private static MemPointIterator mergeSort( @@ -77,12 +89,21 @@ public class MemPointIteratorFactory { List<TimeRange> deletionList, int maxNumberOfPointsInPage) { return new MergeSortMultiTVListIterator( - Ordering.ASC, null, tsDataType, tvLists, deletionList, null, null, maxNumberOfPointsInPage); + Ordering.ASC, + null, + tsDataType, + tvLists, + null, + deletionList, + null, + null, + maxNumberOfPointsInPage); } private static MemPointIterator mergeSort( TSDataType tsDataType, List<TVList> tvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> deletionList, @@ -94,6 +115,7 @@ public class MemPointIteratorFactory { globalTimeFilter, tsDataType, tvLists, + tvListRowCounts, deletionList, floatPrecision, encoding, @@ -104,7 +126,7 @@ public class MemPointIteratorFactory { private static MemPointIterator ordered( TSDataType tsDataType, List<TVList> tvLists, int maxNumberOfPointsInPage) { return new OrderedMultiTVListIterator( - Ordering.ASC, null, tsDataType, tvLists, null, null, null, maxNumberOfPointsInPage); + Ordering.ASC, null, tsDataType, tvLists, null, null, null, null, maxNumberOfPointsInPage); } private static MemPointIterator ordered( @@ -113,12 +135,21 @@ public class MemPointIteratorFactory { List<TimeRange> deletionList, int maxNumberOfPointsInPage) { return new OrderedMultiTVListIterator( - Ordering.ASC, null, tsDataType, tvLists, deletionList, null, null, maxNumberOfPointsInPage); + Ordering.ASC, + null, + tsDataType, + tvLists, + null, + deletionList, + null, + null, + maxNumberOfPointsInPage); } private static MemPointIterator ordered( TSDataType tsDataType, List<TVList> tvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> deletionList, @@ -130,6 +161,7 @@ public class MemPointIteratorFactory { globalTimeFilter, tsDataType, tvLists, + tvListRowCounts, deletionList, floatPrecision, encoding, @@ -147,6 +179,7 @@ public class MemPointIteratorFactory { .get(0) .iterator( Ordering.ASC, + alignedTvLists.get(0).rowCount, null, tsDataTypes, columnIndexList, @@ -170,6 +203,7 @@ public class MemPointIteratorFactory { .get(0) .iterator( Ordering.ASC, + alignedTvLists.get(0).rowCount, null, tsDataTypes, columnIndexList, @@ -185,6 +219,7 @@ public class MemPointIteratorFactory { List<TSDataType> tsDataTypes, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -197,6 +232,7 @@ public class MemPointIteratorFactory { .get(0) .iterator( scanOrder, + tvListRowCounts.get(0), globalTimeFilter, tsDataTypes, columnIndexList, @@ -219,6 +255,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + null, Ordering.ASC, null, null, @@ -241,6 +278,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + null, Ordering.ASC, null, timeColumnDeletion, @@ -255,6 +293,7 @@ public class MemPointIteratorFactory { List<TSDataType> tsDataTypes, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -267,6 +306,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, @@ -288,6 +328,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + null, Ordering.ASC, null, null, @@ -310,6 +351,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + null, Ordering.ASC, null, timeColumnDeletion, @@ -324,6 +366,7 @@ public class MemPointIteratorFactory { List<TSDataType> tsDataTypes, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -336,6 +379,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, @@ -350,7 +394,7 @@ public class MemPointIteratorFactory { TSDataType tsDataType, List<TVList> tvLists, int maxNumberOfPointsInPage) { if (tvLists.size() == 1) { return single(tvLists, maxNumberOfPointsInPage); - } else if (isCompleteOrdered(tvLists)) { + } else if (isCompleteOrdered(tvLists, null)) { return ordered(tsDataType, tvLists, maxNumberOfPointsInPage); } else { return mergeSort(tsDataType, tvLists, maxNumberOfPointsInPage); @@ -364,7 +408,7 @@ public class MemPointIteratorFactory { int maxNumberOfPointsInPage) { if (tvLists.size() == 1) { return single(tvLists, deletionList, maxNumberOfPointsInPage); - } else if (isCompleteOrdered(tvLists)) { + } else if (isCompleteOrdered(tvLists, null)) { return ordered(tsDataType, tvLists, deletionList, maxNumberOfPointsInPage); } else { return mergeSort(tsDataType, tvLists, deletionList, maxNumberOfPointsInPage); @@ -374,6 +418,7 @@ public class MemPointIteratorFactory { public static MemPointIterator create( TSDataType tsDataType, List<TVList> tvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> deletionList, @@ -383,16 +428,18 @@ public class MemPointIteratorFactory { if (tvLists.size() == 1) { return single( tvLists, + tvListRowCounts, scanOrder, globalTimeFilter, deletionList, floatPrecision, encoding, maxNumberOfPointsInPage); - } else if (isCompleteOrdered(tvLists)) { + } else if (isCompleteOrdered(tvLists, tvListRowCounts)) { return ordered( tsDataType, tvLists, + tvListRowCounts, scanOrder, globalTimeFilter, deletionList, @@ -403,6 +450,7 @@ public class MemPointIteratorFactory { return mergeSort( tsDataType, tvLists, + tvListRowCounts, scanOrder, globalTimeFilter, deletionList, @@ -421,7 +469,7 @@ public class MemPointIteratorFactory { if (alignedTvLists.size() == 1) { return single( tsDataTypes, columnIndexList, alignedTvLists, ignoreAllNullRows, maxNumberOfPointsInPage); - } else if (isCompleteOrdered(alignedTvLists)) { + } else if (isCompleteOrdered(alignedTvLists, null)) { return ordered( tsDataTypes, columnIndexList, alignedTvLists, ignoreAllNullRows, maxNumberOfPointsInPage); } else { @@ -447,7 +495,7 @@ public class MemPointIteratorFactory { valueColumnsDeletionList, ignoreAllNullRows, maxNumberOfPointsInPage); - } else if (isCompleteOrdered(alignedTvLists)) { + } else if (isCompleteOrdered(alignedTvLists, null)) { return ordered( tsDataTypes, columnIndexList, @@ -472,6 +520,7 @@ public class MemPointIteratorFactory { List<TSDataType> tsDataTypes, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -485,6 +534,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, @@ -493,11 +543,12 @@ public class MemPointIteratorFactory { encodingList, ignoreAllNullRows, maxNumberOfPointsInPage); - } else if (isCompleteOrdered(alignedTvLists)) { + } else if (isCompleteOrdered(alignedTvLists, tvListRowCounts)) { return ordered( tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, @@ -511,6 +562,7 @@ public class MemPointIteratorFactory { tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, @@ -522,21 +574,24 @@ public class MemPointIteratorFactory { } } - private static boolean isCompleteOrdered(List<? extends TVList> tvLists) { + private static boolean isCompleteOrdered( + List<? extends TVList> tvLists, List<Integer> tvListRowCounts) { long time = Long.MIN_VALUE; for (int i = 0; i < tvLists.size(); i++) { TVList list = tvLists.get(i); - if (!list.isSorted()) { - return false; - } + int rowCount = tvListRowCounts == null ? list.rowCount() : tvListRowCounts.get(i); - if (tvLists.get(i).rowCount() == 0) { + if (rowCount == 0) { continue; } + if (list.seqRowCount() < rowCount) { + return false; + } + if (i > 0 && list.getTime(0) <= time) { return false; } - time = list.getTime(list.rowCount() - 1); + time = list.getTime(rowCount - 1); } return true; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiAlignedTVListIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiAlignedTVListIterator.java index 153cb6856c1..e4f2f3f0d29 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiAlignedTVListIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiAlignedTVListIterator.java @@ -54,6 +54,7 @@ public class MergeSortMultiAlignedTVListIterator extends MultiAlignedTVListItera List<TSDataType> tsDataTypes, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -66,6 +67,7 @@ public class MergeSortMultiAlignedTVListIterator extends MultiAlignedTVListItera tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiTVListIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiTVListIterator.java index d816551a779..f6ed3eebd39 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiTVListIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MergeSortMultiTVListIterator.java @@ -48,6 +48,7 @@ public class MergeSortMultiTVListIterator extends MultiTVListIterator { Filter globalTimeFilter, TSDataType tsDataType, List<TVList> tvLists, + List<Integer> tvListRowCounts, List<TimeRange> deletionList, Integer floatPrecision, TSEncoding encoding, @@ -57,6 +58,7 @@ public class MergeSortMultiTVListIterator extends MultiTVListIterator { globalTimeFilter, tsDataType, tvLists, + tvListRowCounts, deletionList, floatPrecision, encoding, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiAlignedTVListIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiAlignedTVListIterator.java index a2d50dc53a2..31303354b37 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiAlignedTVListIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiAlignedTVListIterator.java @@ -61,6 +61,7 @@ public abstract class MultiAlignedTVListIterator extends MemPointIterator { List<TSDataType> tsDataTypeList, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -74,10 +75,12 @@ public abstract class MultiAlignedTVListIterator extends MemPointIterator { this.columnIndexList = columnIndexList; this.alignedTvListIterators = new ArrayList<>(alignedTvLists.size()); if (scanOrder.isAscending()) { - for (AlignedTVList alignedTVList : alignedTvLists) { + for (int i = 0; i < alignedTvLists.size(); i++) { + AlignedTVList alignedTVList = alignedTvLists.get(i); AlignedTVList.AlignedTVListIterator iterator = alignedTVList.iterator( scanOrder, + tvListRowCounts.get(i), globalTimeFilter, tsDataTypeList, columnIndexList, @@ -95,6 +98,7 @@ public abstract class MultiAlignedTVListIterator extends MemPointIterator { AlignedTVList.AlignedTVListIterator iterator = alignedTVList.iterator( scanOrder, + tvListRowCounts.get(i), globalTimeFilter, tsDataTypeList, columnIndexList, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiTVListIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiTVListIterator.java index a8d5879b767..b735cc7927b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiTVListIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/MultiTVListIterator.java @@ -56,6 +56,7 @@ public abstract class MultiTVListIterator extends MemPointIterator { Filter globalTimeFilter, TSDataType tsDataType, List<TVList> tvLists, + List<Integer> tvListRowCounts, List<TimeRange> deletionList, Integer floatPrecision, TSEncoding encoding, @@ -64,18 +65,33 @@ public abstract class MultiTVListIterator extends MemPointIterator { this.tsDataType = tsDataType; this.tvListIterators = new ArrayList<>(tvLists.size()); if (scanOrder.isAscending()) { - for (TVList tvList : tvLists) { + for (int i = 0; i < tvLists.size(); i++) { + TVList tvList = tvLists.get(i); + int rowCount = tvListRowCounts == null ? tvList.rowCount : tvListRowCounts.get(i); TVList.TVListIterator iterator = tvList.iterator( - scanOrder, globalTimeFilter, deletionList, null, null, maxNumberOfPointsInPage); + scanOrder, + rowCount, + globalTimeFilter, + deletionList, + null, + null, + maxNumberOfPointsInPage); tvListIterators.add(iterator); } } else { for (int i = tvLists.size() - 1; i >= 0; i--) { TVList tvList = tvLists.get(i); + int rowCount = tvListRowCounts == null ? tvList.rowCount : tvListRowCounts.get(i); TVList.TVListIterator iterator = tvList.iterator( - scanOrder, globalTimeFilter, deletionList, null, null, maxNumberOfPointsInPage); + scanOrder, + rowCount, + globalTimeFilter, + deletionList, + null, + null, + maxNumberOfPointsInPage); tvListIterators.add(iterator); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiAlignedTVListIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiAlignedTVListIterator.java index 1c405f01bfe..4a167f174a3 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiAlignedTVListIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiAlignedTVListIterator.java @@ -41,6 +41,7 @@ public class OrderedMultiAlignedTVListIterator extends MultiAlignedTVListIterato List<TSDataType> tsDataTypes, List<Integer> columnIndexList, List<AlignedTVList> alignedTvLists, + List<Integer> tvListRowCounts, Ordering scanOrder, Filter globalTimeFilter, List<TimeRange> timeColumnDeletion, @@ -53,6 +54,7 @@ public class OrderedMultiAlignedTVListIterator extends MultiAlignedTVListIterato tsDataTypes, columnIndexList, alignedTvLists, + tvListRowCounts, scanOrder, globalTimeFilter, timeColumnDeletion, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiTVListIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiTVListIterator.java index c7ae99386c5..61ae5f1360d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiTVListIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/OrderedMultiTVListIterator.java @@ -35,6 +35,7 @@ public class OrderedMultiTVListIterator extends MultiTVListIterator { Filter globalTimeFilter, TSDataType tsDataType, List<TVList> tvLists, + List<Integer> tvListRowCounts, List<TimeRange> deletionList, Integer floatPrecision, TSEncoding encoding, @@ -44,6 +45,7 @@ public class OrderedMultiTVListIterator extends MultiTVListIterator { globalTimeFilter, tsDataType, tvLists, + tvListRowCounts, deletionList, floatPrecision, encoding, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java index dedad0ee793..4309d1e68f2 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java @@ -656,6 +656,7 @@ public abstract class TVList implements WALEntryValue { public TVListIterator iterator( Ordering scanOrder, + int rowCount, Filter globalTimeFilter, List<TimeRange> deletionList, Integer floatPrecision, @@ -663,6 +664,7 @@ public abstract class TVList implements WALEntryValue { int maxNumberOfPointsInPage) { return new TVListIterator( scanOrder, + rowCount, globalTimeFilter, deletionList, floatPrecision, @@ -687,6 +689,7 @@ public abstract class TVList implements WALEntryValue { public TVListIterator( Ordering scanOrder, + int rowCount, Filter globalTimeFilter, List<TimeRange> deletionList, Integer floatPrecision, diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedTVListIteratorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedTVListIteratorTest.java index f0ae227c05f..f02a637dad5 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedTVListIteratorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedTVListIteratorTest.java @@ -774,7 +774,6 @@ public class AlignedTVListIteratorTest { AlignedTVList alignedTVList = AlignedTVList.newAlignedList( Arrays.asList(TSDataType.INT64, TSDataType.BOOLEAN, TSDataType.BOOLEAN)); - int rowCount = 0; for (TimeRange timeRange : timeRanges) { long start = timeRange.getMin(); long end = timeRange.getMax(); @@ -793,11 +792,10 @@ public class AlignedTVListIteratorTest { } alignedTVList.putAlignedValue( timestamp, new Object[] {timestamp, timestamp % 2 == 0, true}); - rowCount++; } } Map<TVList, Integer> tvListMap = new HashMap<>(); - tvListMap.put(alignedTVList, rowCount); + tvListMap.put(alignedTVList, alignedTVList.rowCount()); return tvListMap; } @@ -808,7 +806,6 @@ public class AlignedTVListIteratorTest { AlignedTVList alignedTVList = AlignedTVList.newAlignedList( Arrays.asList(TSDataType.INT64, TSDataType.BOOLEAN, TSDataType.BOOLEAN)); - int rowCount = 0; long start = timeRange.getMin(); long end = timeRange.getMax(); List<Long> timestamps = new ArrayList<>((int) (end - start + 1)); @@ -826,9 +823,8 @@ public class AlignedTVListIteratorTest { } alignedTVList.putAlignedValue( timestamp, new Object[] {timestamp, timestamp % 2 == 0, isLast}); - rowCount++; } - tvListMap.put(alignedTVList, rowCount); + tvListMap.put(alignedTVList, alignedTVList.rowCount()); } return tvListMap; } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/NonAlignedTVListIteratorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/NonAlignedTVListIteratorTest.java index 62c801f4e34..a984eaada1f 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/NonAlignedTVListIteratorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/NonAlignedTVListIteratorTest.java @@ -582,7 +582,6 @@ public class NonAlignedTVListIteratorTest { private static Map<TVList, Integer> buildNonAlignedSingleTvListMap(List<TimeRange> timeRanges) { TVList tvList = TVList.newList(TSDataType.INT64); - int rowCount = 0; for (TimeRange timeRange : timeRanges) { long start = timeRange.getMin(); long end = timeRange.getMax(); @@ -599,11 +598,10 @@ public class NonAlignedTVListIteratorTest { } } tvList.putLong(timestamp, timestamp); - rowCount++; } } Map<TVList, Integer> tvListMap = new HashMap<>(); - tvListMap.put(tvList, rowCount); + tvListMap.put(tvList, tvList.rowCount()); return tvListMap; } @@ -611,7 +609,6 @@ public class NonAlignedTVListIteratorTest { Map<TVList, Integer> tvListMap = new LinkedHashMap<>(); for (TimeRange timeRange : timeRanges) { TVList tvList = TVList.newList(TSDataType.INT64); - int rowCount = 0; long start = timeRange.getMin(); long end = timeRange.getMax(); List<Long> timestamps = new ArrayList<>((int) (end - start + 1)); @@ -621,9 +618,8 @@ public class NonAlignedTVListIteratorTest { Collections.shuffle(timestamps); for (Long timestamp : timestamps) { tvList.putLong(timestamp, timestamp); - rowCount++; } - tvListMap.put(tvList, rowCount); + tvListMap.put(tvList, tvList.rowCount()); } return tvListMap; }
