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..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 @@ -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() { @@ -1238,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/AlignedTVListIteratorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedTVListIteratorTest.java index d5a8b49f7260c..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 @@ -846,6 +846,90 @@ 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, + 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, + 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 = 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..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 @@ -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,73 @@ 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); + + 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); + + 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()); + } }