Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -966,9 +966,10 @@ protected void skipToCurrentTimeRangeStartPosition() {
this.getQueryContext().getQueryStatistics().addFilteredRowsOfRowLevel(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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Invalidate non-aligned batch state before preserving probeNext

newIndex == index does not guarantee that probeNext still describes the current record. Non-aligned TVListIterator.nextBatch() advances index directly but does not clear probeNext before returning. With this change, switching through an empty fake page preserves that stale true, so subsequent point reads skip prepareNext() and can return duplicate timestamps or deleted points. This matters when an earlier non-overlapping page is read in batches and a later overlapping page is consumed through the point-reader path.

For example, consider a single non-aligned LongTVList with these physical records, in write order:

index  time  value
0      1     1
1      100   2
2      100   3

Read [1,33] in batches, visit the empty page [34,66], then read [67,100] as points:

  1. The first hasNextBatch() prepares index 0 and sets probeNext = true.
  2. nextBatch() emits (1,1), advances to index 1, and stops because time 100 exceeds the first page boundary. It leaves probeNext = true, although index 1 has not been prepared for point reading.
  3. For the empty page, binary search returns index 1 again. Previously, this branch cleared probeNext, allowing hasNextBatch() to prepare the time-100 group and advance to its latest record, index 2. After this change, preparation is skipped and index 1 remains selected.
  4. When the last page is selected, time 100 is already inside its range, so the early return leaves the stale flag untouched. The first point read emits (100,2); the next read also emits (100,3).
Expected last-page output: [(100,3)]
Parent implementation:    [(100,3)]
This PR:                  [(100,2), (100,3)]

The same stale flag bypasses deletion checks. My ASC and DESC deletion variants both passed with the parent implementation and returned an extra deleted point with this PR.

Please clear probeNext after non-aligned nextBatch() advances the cursor, as the aligned batch implementation already does, and cover the batch-to-point transition while retaining these aligned regressions.

Minimal JUnit reproducer (imports omitted)
@Test
public void testNonAlignedBatchToPointAfterEmptyPage() throws Exception {
  LongTVList list = LongTVList.newList();
  list.putLong(1, 1);
  list.putLong(100, 2);
  list.putLong(100, 3);

  MemPointIterator iterator =
      list.iterator(
          Ordering.ASC,
          list.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<String> result = new ArrayList<>();
  while (iterator.hasNextTimeValuePair()) {
    TimeValuePair pair = iterator.nextTimeValuePair();
    result.add(pair.getTimestamp() + ":" + pair.getValue().getLong());
  }
  Assert.assertEquals(Collections.singletonList("100:3"), result);
}

This minimal test supplies the page ranges explicitly. I also reproduced it with 300,000 physical rows: 299,998 records at time 1 followed by (100,2) and (100,3). ReadOnlyMemChunk naturally generated [1,33], [34,66], and [67,100]; the same batch/batch/point sequence through LazyMemVersionPageReader and PriorityMergeReader returned [100:2, 100:3] on this PR and [100:3] with the parent implementation.

}

probeNext = false;
}

protected void prepareNext() {
Expand Down Expand Up @@ -1312,6 +1313,7 @@ public TsBlock nextBatch() {
// been applied when constructing the tsBlock
TsBlock tsBlock = builder.build();
addTsBlock(tsBlock);
probeNext = false;
return tsBlock;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -868,6 +868,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<TVList, Integer> 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<TimeValuePair> 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<TVList, Integer> 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<TimeValuePair> 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<Map<TVList, Integer>> list =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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<Long> 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<Long> result = new ArrayList<>();
while (iterator.hasNextTimeValuePair()) {
result.add(iterator.nextTimeValuePair().getValue().getLong());
}
Assert.assertTrue(result.isEmpty());
}
}
Loading