Skip to content

Commit 5f0b506

Browse files
committed
[Subscription] Keep WAL iterator across bounded gap rounds
1 parent b777d3c commit 5f0b506

2 files changed

Lines changed: 27 additions & 12 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1861,10 +1861,15 @@ private MaterializationResult fillGapFromWAL(
18611861
final int maxTablets,
18621862
final long maxBatchBytes) {
18631863
pendingWalGapRetryRequested = false;
1864-
resetSubscriptionWALPosition(fromIndex);
18651864
if (seekGeneration.get() != expectedSeekGeneration || isClosed) {
18661865
return MaterializationResult.STALE;
18671866
}
1867+
// Rebuilding the iterator for every bounded gap round rereads the retained WAL prefix to
1868+
// reach fromIndex. Keep its current position; advanceTo skips whole files only when their
1869+
// writer progress is already covered, and otherwise replays them sequentially.
1870+
if (Objects.nonNull(subscriptionWALIterator)) {
1871+
subscriptionWALIterator.advanceTo(fromIndex, this::isWriterProgressCoveredForWalFastForward);
1872+
}
18681873
// A pending-queue jump can span millions of WAL entries when real-time delivery has overflowed.
18691874
// Keep gap recovery within the normal per-round budget so sparse topic filtering cannot hold
18701875
// the

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java‎

Lines changed: 21 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1410,17 +1410,8 @@ public void testPendingGapReplayHonorsPerRoundWalEntryLimit() throws Exception {
14101410
});
14111411
when(converter.getDatabaseName()).thenReturn("db");
14121412

1413-
final Iterator<IndexedConsensusRequest> walEntries =
1414-
Arrays.asList(
1415-
createRequest(1L),
1416-
createRequest(2L),
1417-
createRequest(3L),
1418-
createRequest(4L),
1419-
createRequest(5L))
1420-
.iterator();
1421-
final ProgressWALIterator iterator = mock(ProgressWALIterator.class);
1422-
when(iterator.hasNext()).thenAnswer(ignored -> walEntries.hasNext());
1423-
when(iterator.next()).thenAnswer(ignored -> walEntries.next());
1413+
final AtomicInteger iteratorCreationCount = new AtomicInteger();
1414+
final AtomicInteger walEntryReadCount = new AtomicInteger();
14241415

14251416
queue =
14261417
new ConsensusPrefetchingQueue(
@@ -1442,6 +1433,23 @@ public void testPendingGapReplayHonorsPerRoundWalEntryLimit() throws Exception {
14421433
@Override
14431434
protected ProgressWALIterator createSubscriptionWALIterator(
14441435
final long startSearchIndex) {
1436+
iteratorCreationCount.incrementAndGet();
1437+
final Iterator<IndexedConsensusRequest> walEntries =
1438+
Arrays.asList(
1439+
createRequest(1L),
1440+
createRequest(2L),
1441+
createRequest(3L),
1442+
createRequest(4L),
1443+
createRequest(5L))
1444+
.iterator();
1445+
final ProgressWALIterator iterator = mock(ProgressWALIterator.class);
1446+
when(iterator.hasNext()).thenAnswer(ignored -> walEntries.hasNext());
1447+
when(iterator.next())
1448+
.thenAnswer(
1449+
ignored -> {
1450+
walEntryReadCount.incrementAndGet();
1451+
return walEntries.next();
1452+
});
14451453
return iterator;
14461454
}
14471455
};
@@ -1462,6 +1470,8 @@ protected ProgressWALIterator createSubscriptionWALIterator(
14621470
queue.drivePrefetchOnce();
14631471
assertEquals(5, conversionCount.get());
14641472
assertEquals(6L, queue.getCurrentReadSearchIndex());
1473+
assertEquals(1, iteratorCreationCount.get());
1474+
assertEquals(5, walEntryReadCount.get());
14651475
} finally {
14661476
if (queue != null) {
14671477
queue.close();

0 commit comments

Comments
 (0)