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 @@ -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() {
Expand Down Expand Up @@ -1238,6 +1239,7 @@ && isTimeSatisfied(time)) {
// 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 @@ -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<TVList, Integer> 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<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,
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,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<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);

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