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;
   }

Reply via email to