This is an automated email from the ASF dual-hosted git repository.
jackietien pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 7466e2ce1f0 Concurrently querying and writing to the memtable may
cause the query results out of order (#16328)
7466e2ce1f0 is described below
commit 7466e2ce1f0a083649d7669f68f00e720e017b05
Author: shuwenwei <[email protected]>
AuthorDate: Fri Sep 5 09:21:08 2025 +0800
Concurrently querying and writing to the memtable may cause the query
results out of order (#16328)
---
.../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 | 9 +--
.../memtable/AlignedTVListIteratorTest.java | 8 +-
.../memtable/NonAlignedTVListIteratorTest.java | 8 +-
13 files changed, 134 insertions(+), 75 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..42de51ac4b7 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 == null ? alignedTVList.rowCount() :
tvListRowCounts.get(i),
globalTimeFilter,
tsDataTypeList,
columnIndexList,
@@ -95,6 +98,7 @@ public abstract class MultiAlignedTVListIterator extends
MemPointIterator {
AlignedTVList.AlignedTVListIterator iterator =
alignedTVList.iterator(
scanOrder,
+ tvListRowCounts == null ? alignedTVList.rowCount() :
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 06d27fcf2cf..00b49954470 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
@@ -251,10 +251,6 @@ public abstract class TVList implements WALEntryValue {
return indices.get(arrayIndex)[elementIndex];
}
- public int getValueIndex(int index, Ordering ordering) {
- return ordering.isAscending() ? getValueIndex(index) :
getValueIndex(rowCount - 1 - index);
- }
-
protected void markNullValue(int arrayIndex, int elementIndex) {
// init bitMap if doesn't have
if (bitMap == null) {
@@ -656,6 +652,7 @@ public abstract class TVList implements WALEntryValue {
public TVListIterator iterator(
Ordering scanOrder,
+ int rowCount,
Filter globalTimeFilter,
List<TimeRange> deletionList,
Integer floatPrecision,
@@ -663,6 +660,7 @@ public abstract class TVList implements WALEntryValue {
int maxNumberOfPointsInPage) {
return new TVListIterator(
scanOrder,
+ rowCount,
globalTimeFilter,
deletionList,
floatPrecision,
@@ -687,6 +685,7 @@ public abstract class TVList implements WALEntryValue {
public TVListIterator(
Ordering scanOrder,
+ int rowCount,
Filter globalTimeFilter,
List<TimeRange> deletionList,
Integer floatPrecision,
@@ -966,7 +965,7 @@ public abstract class TVList implements WALEntryValue {
// When traversing in desc order, the index needs to be converted
public int getScanOrderIndex(int rowIndex) {
- return scanOrder.isAscending() ? rowIndex : rowCount - 1 - rowIndex;
+ return scanOrder.isAscending() ? rowIndex : rows - 1 - rowIndex;
}
@Override
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 6da636e47e5..2a3cabdf08c 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
@@ -780,7 +780,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();
@@ -799,11 +798,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;
}
@@ -814,7 +812,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));
@@ -832,9 +829,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 2b06382161f..c22dcbb6b0f 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
@@ -589,7 +589,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();
@@ -606,11 +605,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;
}
@@ -618,7 +616,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));
@@ -628,9 +625,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;
}