Skip to content

Commit bf8ac51

Browse files
authored
[Subscription] Reduce consensus WAL replay and ACK contention (#18402)
1 parent ab31e1b commit bf8ac51

4 files changed

Lines changed: 895 additions & 57 deletions

File tree

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

Lines changed: 57 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -193,9 +193,9 @@ public class ConsensusPrefetchingQueue {
193193
private volatile ProgressWALIterator subscriptionWALIterator;
194194

195195
/**
196-
* Seek requests must not close/reset the WAL iterator from RPC threads because the prefetch
197-
* worker may be reading it concurrently. Instead, seek only records the latest desired reset and
198-
* the queue's next prefetch round applies it after observing the new seek generation.
196+
* WAL cursor changes outside the iterator must not close/reset it from RPC threads because the
197+
* prefetch worker may be reading it concurrently. Instead, the latest desired reset is recorded
198+
* and applied by the next prefetch round after observing the expected seek generation.
199199
*/
200200
private volatile long pendingSubscriptionWalResetSearchIndex = Long.MIN_VALUE;
201201

@@ -1567,12 +1567,41 @@ private boolean isBeforeLocalCursor(final IndexedConsensusRequest request) {
15671567
return hasLocalSearchIndex(request) && request.getSearchIndex() < nextExpectedSearchIndex.get();
15681568
}
15691569

1570-
private void advanceLocalCursorIfPresent(final IndexedConsensusRequest request) {
1570+
private boolean advanceLocalCursorIfPresent(final IndexedConsensusRequest request) {
15711571
if (hasLocalSearchIndex(request)) {
15721572
nextExpectedSearchIndex.set(request.getSearchIndex() + 1);
1573+
return true;
1574+
}
1575+
return false;
1576+
}
1577+
1578+
private void advanceLocalCursorFromPendingIfPresent(
1579+
final IndexedConsensusRequest request, final long expectedSeekGeneration) {
1580+
if (advanceLocalCursorIfPresent(request)) {
1581+
// Pending delivery advances independently of the WAL reader. Raise its local lower bound in
1582+
// place so stale local requests are filtered without rebuilding and rescanning retained WAL.
1583+
final ProgressWALIterator iterator = subscriptionWALIterator;
1584+
if (Objects.nonNull(iterator) && seekGeneration.get() == expectedSeekGeneration) {
1585+
iterator.advanceTo(
1586+
nextExpectedSearchIndex.get(), this::isWriterProgressCoveredForWalFastForward);
1587+
}
15731588
}
15741589
}
15751590

1591+
private boolean isWriterProgressCoveredForWalFastForward(
1592+
final long physicalTime, final int nodeId, final long localSeq) {
1593+
final WriterProgress candidate = new WriterProgress(physicalTime, localSeq);
1594+
final WriterProgress recoveryProgress =
1595+
recoveryWriterProgressByWriter.get(new WriterId(consensusGroupId.toString(), nodeId));
1596+
if (Objects.nonNull(recoveryProgress)
1597+
&& compareWriterProgress(candidate, recoveryProgress) <= 0) {
1598+
return true;
1599+
}
1600+
final WriterProgress materializedProgress = materializedProgressByWriter.get(nodeId);
1601+
return Objects.nonNull(materializedProgress)
1602+
&& compareWriterProgress(candidate, materializedProgress) <= 0;
1603+
}
1604+
15761605
private MaterializationResult appendRealtimeRequest(
15771606
final IndexedConsensusRequest request,
15781607
final DeliveryBatchState batchState,
@@ -1648,12 +1677,12 @@ private MaterializationResult accumulateFromPending(
16481677

16491678
if (shouldSkipForRecoveryProgress(request)) {
16501679
skippedCount++;
1651-
advanceLocalCursorIfPresent(request);
1680+
advanceLocalCursorFromPendingIfPresent(request, expectedSeekGeneration);
16521681
continue;
16531682
}
16541683
if (shouldSkipForMaterializedProgress(request)) {
16551684
skippedCount++;
1656-
advanceLocalCursorIfPresent(request);
1685+
advanceLocalCursorFromPendingIfPresent(request, expectedSeekGeneration);
16571686
continue;
16581687
}
16591688

@@ -1665,7 +1694,7 @@ private MaterializationResult accumulateFromPending(
16651694
}
16661695
markMaterializedProgress(request);
16671696
processedCount++;
1668-
advanceLocalCursorIfPresent(request);
1697+
advanceLocalCursorFromPendingIfPresent(request, expectedSeekGeneration);
16691698
if (prefetchingQueue.size() >= MAX_PREFETCHING_QUEUE_SIZE) {
16701699
break;
16711700
}
@@ -1757,7 +1786,9 @@ private MaterializationResult tryCatchUpFromWAL(final long expectedSeekGeneratio
17571786
// Use the persistent linger batch so an unexpected runtime failure cannot orphan already
17581787
// reserved Tablets or advance replay progress past data that has become unreachable.
17591788
final DeliveryBatchState batchState = lingerBatch;
1760-
resetSubscriptionWALPosition(nextExpectedSearchIndex.get());
1789+
// Keep the iterator and its buffered next request across rounds. Reopening it here discards the
1790+
// request prepared by hasNext() and repeatedly re-reads, skips, and decompresses the same WAL
1791+
// segment. Pending-path cursor advances and seek operations request explicit realignment.
17611792
final MaterializationResult materializationResult =
17621793
pumpFromSubscriptionWAL(
17631794
batchState, expectedSeekGeneration, maxWalEntries, maxTablets, maxBatchBytes);
@@ -1790,7 +1821,6 @@ private MaterializationResult pumpFromSubscriptionWAL(
17901821
return MaterializationResult.SUCCESS;
17911822
}
17921823

1793-
subscriptionWALIterator.refresh();
17941824
ensureSubscriptionWalReadable();
17951825

17961826
int entriesRead = 0;
@@ -1846,9 +1876,14 @@ private MaterializationResult pumpFromSubscriptionWAL(
18461876
}
18471877

18481878
private void ensureSubscriptionWalReadable() {
1849-
if (Objects.isNull(subscriptionWALIterator)
1850-
|| subscriptionWALIterator.hasNext()
1851-
|| !(consensusReqReader instanceof WALNode)) {
1879+
if (Objects.isNull(subscriptionWALIterator) || subscriptionWALIterator.hasNext()) {
1880+
return;
1881+
}
1882+
1883+
// Listing and sorting all retained WAL files is only necessary after the iterator is
1884+
// exhausted. While it still has a readable request, refreshing cannot affect the next result.
1885+
subscriptionWALIterator.refresh();
1886+
if (subscriptionWALIterator.hasNext() || !(consensusReqReader instanceof WALNode)) {
18521887
return;
18531888
}
18541889

@@ -1865,9 +1900,6 @@ private void ensureSubscriptionWalReadable() {
18651900
currentWalIndex);
18661901
((WALNode) consensusReqReader).rollWALFile();
18671902
resetSubscriptionWALPosition(nextExpectedSearchIndex.get());
1868-
if (Objects.nonNull(subscriptionWALIterator)) {
1869-
subscriptionWALIterator.refresh();
1870-
}
18711903
}
18721904

18731905
private void resetSubscriptionWALPosition(final long startSearchIndex) {
@@ -1885,6 +1917,11 @@ protected ProgressWALIterator createSubscriptionWALIterator(final long startSear
18851917
protected void onWalGapRetryScheduled() {}
18861918

18871919
private boolean hasReadableWalEntries() {
1920+
if (pendingSubscriptionWalResetSearchIndex != Long.MIN_VALUE) {
1921+
// Do not advance the stale iterator only to discard its buffered request when the next round
1922+
// applies the pending realignment. Returning true keeps the worker scheduled for that round.
1923+
return true;
1924+
}
18881925
return Objects.nonNull(subscriptionWALIterator) && subscriptionWALIterator.hasNext();
18891926
}
18901927

@@ -2179,7 +2216,10 @@ private void cleanUpEvent(final SubscriptionEvent event, final boolean force) {
21792216

21802217
private boolean ackMissingInFlightEvent(
21812218
final SubscriptionCommitContext commitContext, final boolean silent) {
2182-
acquireWriteLock();
2219+
// Late or duplicate ACKs touch the same concurrent lifecycle indexes and commit manager as the
2220+
// regular in-flight ACK path. A read lock is sufficient to fence seek/close transitions while
2221+
// allowing ACKs to proceed concurrently with a long-running WAL prefetch round.
2222+
acquireReadLock();
21832223
try {
21842224
if (!canAcceptCommitContext(commitContext, "ack", silent)) {
21852225
return false;
@@ -2226,7 +2266,7 @@ private boolean ackMissingInFlightEvent(
22262266
}
22272267
return true;
22282268
} finally {
2229-
releaseWriteLock();
2269+
releaseReadLock();
22302270
}
22312271
}
22322272

0 commit comments

Comments
 (0)