From 839e0ca3587863cee77674d31d5dc8aa9edebbb9 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Wed, 16 Sep 2026 12:00:45 +0800 Subject: [PATCH 1/4] Fix aligned TVList page switch dropping duplicate values --- .../iotdb/db/utils/datastructure/TVList.java | 5 +- .../memtable/AlignedTVListIteratorTest.java | 86 +++++++++++++++++++ 2 files changed, 89 insertions(+), 2 deletions(-) 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 2536836069ca5..e0d8e5b6de340 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 @@ -931,9 +931,10 @@ protected void skipToCurrentTimeRangeStartPosition() { int newIndex = getScanOrderIndex(indexInTVList); if (newIndex > index) { index = newIndex; + // If the cursor does not move, a duplicate-timestamp group prepared for the current + // position remains valid. Invalidate it only after the cursor actually advances. + probeNext = false; } - - probeNext = false; } protected void prepareNext() { 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 d5a8b49f7260c..840e9d3115667 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 @@ -846,6 +846,92 @@ private PaginationController duplicatePaginationController( paginationController.getCurLimit(), paginationController.getCurOffset()); } + @Test + public void testPageSwitchKeepsPreparedDuplicateTimestampValues() throws IOException { + AlignedTVList tvList = + AlignedTVList.newAlignedList( + Arrays.asList(TSDataType.INT64, TSDataType.BOOLEAN, TSDataType.BOOLEAN)); + tvList.putAlignedValue(1, new Object[] {1L, true, false}); + tvList.putAlignedValue(100, new Object[] {2L, null, false}); + tvList.putAlignedValue(100, new Object[] {null, true, false}); + + Map tvListMap = new LinkedHashMap<>(); + tvListMap.put(tvList, tvList.rowCount()); + AlignedReadOnlyMemChunk chunk = + new AlignedReadOnlyMemChunk( + fragmentInstanceContext, + Arrays.asList(0, 1, 2), + getMeasurementSchema(), + tvListMap, + Collections.emptyList(), + Arrays.asList( + Collections.emptyList(), Collections.emptyList(), Collections.emptyList())); + chunk.sortTvLists(); + chunk.initChunkMetaFromTVListsWithFakeStatistics(); + + MemPointIterator iterator = chunk.createMemPointIterator(Ordering.ASC, null); + List result = new ArrayList<>(); + // These are fake-page boundaries for one MemChunk. The middle page is empty, but the + // shared iterator still receives its time range before its next page is read. + for (TimeRange pageRange : + Arrays.asList(new TimeRange(1, 33), new TimeRange(34, 66), new TimeRange(67, 100))) { + iterator.setCurrentPageTimeRange(pageRange); + while (iterator.hasNextTimeValuePair()) { + result.add(iterator.nextTimeValuePair()); + } + } + + Assert.assertEquals(2, result.size()); + Assert.assertEquals(1L, result.get(0).getTimestamp()); + Assert.assertEquals(1L, result.get(0).getValues()[0]); + Assert.assertEquals(100L, result.get(1).getTimestamp()); + Assert.assertEquals(2L, result.get(1).getValues()[0]); + Assert.assertEquals(Boolean.TRUE, result.get(1).getValues()[1]); + Assert.assertEquals(Boolean.FALSE, result.get(1).getValues()[2]); + } + + @Test + public void testPageSwitchKeepsPreparedDuplicateTimestampValuesDescending() throws IOException { + AlignedTVList tvList = + AlignedTVList.newAlignedList( + Arrays.asList(TSDataType.INT64, TSDataType.BOOLEAN, TSDataType.BOOLEAN)); + tvList.putAlignedValue(1, new Object[] {null, true, false}); + tvList.putAlignedValue(1, new Object[] {2L, null, false}); + tvList.putAlignedValue(100, new Object[] {1L, true, false}); + + Map tvListMap = new LinkedHashMap<>(); + tvListMap.put(tvList, tvList.rowCount()); + AlignedReadOnlyMemChunk chunk = + new AlignedReadOnlyMemChunk( + fragmentInstanceContext, + Arrays.asList(0, 1, 2), + getMeasurementSchema(), + tvListMap, + Collections.emptyList(), + Arrays.asList( + Collections.emptyList(), Collections.emptyList(), Collections.emptyList())); + chunk.sortTvLists(); + chunk.initChunkMetaFromTVListsWithFakeStatistics(); + + MemPointIterator iterator = chunk.createMemPointIterator(Ordering.DESC, null); + List result = new ArrayList<>(); + for (TimeRange pageRange : + Arrays.asList(new TimeRange(67, 100), new TimeRange(34, 66), new TimeRange(1, 33))) { + iterator.setCurrentPageTimeRange(pageRange); + while (iterator.hasNextTimeValuePair()) { + result.add(iterator.nextTimeValuePair()); + } + } + + Assert.assertEquals(2, result.size()); + Assert.assertEquals(100L, result.get(0).getTimestamp()); + Assert.assertEquals(1L, result.get(0).getValues()[0]); + Assert.assertEquals(1L, result.get(1).getTimestamp()); + Assert.assertEquals(2L, result.get(1).getValues()[0]); + Assert.assertEquals(Boolean.TRUE, result.get(1).getValues()[1]); + Assert.assertEquals(Boolean.FALSE, result.get(1).getValues()[2]); + } + @Test public void testSkipTimeRange() throws QueryProcessException, IOException { List> list = From 38dd721c7f7c47083f63a447cfc93434aaeb7042 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Wed, 16 Sep 2026 15:08:23 +0800 Subject: [PATCH 2/4] Fix non-aligned TVList iterator stale prepared state --- .../iotdb/db/utils/datastructure/TVList.java | 1 + .../NonAlignedTVListIteratorTest.java | 72 +++++++++++++++++++ 2 files changed, 73 insertions(+) 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 e0d8e5b6de340..4826021314054 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 @@ -1239,6 +1239,7 @@ && isTimeSatisfied(time)) { // been applied when constructing the tsBlock TsBlock tsBlock = builder.build(); addTsBlock(tsBlock); + probeNext = false; return tsBlock; } 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 b28979efd677a..514fb93db4fa9 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 @@ -26,6 +26,7 @@ import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; import org.apache.iotdb.db.queryengine.plan.statement.component.Ordering; +import org.apache.iotdb.db.utils.datastructure.LongTVList; import org.apache.iotdb.db.utils.datastructure.MemPointIterator; import org.apache.iotdb.db.utils.datastructure.TVList; @@ -717,4 +718,75 @@ private void testSkipTimeRange( } Assert.assertEquals(expectedTimestamps, resultTimestamps); } + + @Test + public void testBatchToPointAfterEmptyPageKeepsLatestDuplicateValue() throws IOException { + LongTVList tvList = LongTVList.newList(); + tvList.putLong(1, 1); + tvList.putLong(100, 2); + tvList.putLong(100, 3); + + MemPointIterator iterator = + tvList.iterator( + Ordering.ASC, + tvList.rowCount(), + null, + Collections.emptyList(), + 0, + TSEncoding.PLAIN, + 1024, + null); + + iterator.setCurrentPageTimeRange(new TimeRange(1, 33)); + int firstPageRows = 0; + while (iterator.hasNextBatch()) { + firstPageRows += iterator.nextBatch().getPositionCount(); + } + Assert.assertEquals(1, firstPageRows); + + iterator.setCurrentPageTimeRange(new TimeRange(34, 66)); + Assert.assertFalse(iterator.hasNextBatch()); + + iterator.setCurrentPageTimeRange(new TimeRange(67, 100)); + List result = new ArrayList<>(); + while (iterator.hasNextTimeValuePair()) { + result.add(iterator.nextTimeValuePair().getValue().getLong()); + } + Assert.assertEquals(Collections.singletonList(3L), result); + } + + @Test + public void testBatchToPointAfterEmptyPageDescendingSkipsDeletedPoint() throws IOException { + LongTVList tvList = LongTVList.newList(); + tvList.putLong(10, 10); + tvList.putLong(100, 100); + + MemPointIterator iterator = + tvList.iterator( + Ordering.DESC, + tvList.rowCount(), + null, + Collections.singletonList(new TimeRange(10, 10)), + 0, + TSEncoding.PLAIN, + 1024, + null); + + iterator.setCurrentPageTimeRange(new TimeRange(67, 100)); + int firstPageRows = 0; + while (iterator.hasNextBatch()) { + firstPageRows += iterator.nextBatch().getPositionCount(); + } + Assert.assertEquals(1, firstPageRows); + + iterator.setCurrentPageTimeRange(new TimeRange(34, 66)); + Assert.assertFalse(iterator.hasNextBatch()); + + iterator.setCurrentPageTimeRange(new TimeRange(1, 33)); + List result = new ArrayList<>(); + while (iterator.hasNextTimeValuePair()) { + result.add(iterator.nextTimeValuePair().getValue().getLong()); + } + Assert.assertTrue(result.isEmpty()); + } } From 2eae01522b2ba5b5bfed1981cff0bcb127239df4 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Wed, 16 Sep 2026 16:28:03 +0800 Subject: [PATCH 3/4] Adapt non-aligned iterator tests to dev/1.3 API --- .../dataregion/memtable/NonAlignedTVListIteratorTest.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) 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 514fb93db4fa9..3f63812447a5f 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 @@ -734,8 +734,7 @@ public void testBatchToPointAfterEmptyPageKeepsLatestDuplicateValue() throws IOE Collections.emptyList(), 0, TSEncoding.PLAIN, - 1024, - null); + 1024); iterator.setCurrentPageTimeRange(new TimeRange(1, 33)); int firstPageRows = 0; @@ -769,8 +768,7 @@ public void testBatchToPointAfterEmptyPageDescendingSkipsDeletedPoint() throws I Collections.singletonList(new TimeRange(10, 10)), 0, TSEncoding.PLAIN, - 1024, - null); + 1024); iterator.setCurrentPageTimeRange(new TimeRange(67, 100)); int firstPageRows = 0; From 32e5db6181a143d6b855ec1d81c72d8635cdeb30 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Wed, 16 Sep 2026 16:29:30 +0800 Subject: [PATCH 4/4] Adapt aligned iterator tests to dev/1.3 API --- .../dataregion/memtable/AlignedTVListIteratorTest.java | 2 -- 1 file changed, 2 deletions(-) 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 840e9d3115667..04ea5958a2bf6 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 @@ -863,7 +863,6 @@ public void testPageSwitchKeepsPreparedDuplicateTimestampValues() throws IOExcep Arrays.asList(0, 1, 2), getMeasurementSchema(), tvListMap, - Collections.emptyList(), Arrays.asList( Collections.emptyList(), Collections.emptyList(), Collections.emptyList())); chunk.sortTvLists(); @@ -907,7 +906,6 @@ public void testPageSwitchKeepsPreparedDuplicateTimestampValuesDescending() thro Arrays.asList(0, 1, 2), getMeasurementSchema(), tvListMap, - Collections.emptyList(), Arrays.asList( Collections.emptyList(), Collections.emptyList(), Collections.emptyList())); chunk.sortTvLists();