From c3825538270c102d850be0c9c36c2a0ca19d840f Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Mon, 28 Sep 2026 16:25:18 +0800 Subject: [PATCH 1/4] Fix silent subscription data loss on WAL replay gaps --- .../iotdb/db/i18n/DataNodePipeMessages.java | 11 +- .../iotdb/db/i18n/DataNodePipeMessages.java | 9 +- .../consensus/ConsensusPrefetchingQueue.java | 109 ++++++++++++--- .../ConsensusPrefetchingQueueTest.java | 131 ++++++++++++++++-- 4 files changed, 227 insertions(+), 33 deletions(-) diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 3f7330a79ad8..c9a146d1dd6e 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2226,10 +2226,13 @@ private DataNodePipeMessages() {} "ProgressWALIterator: skipped {} unreadable retained WAL files in directory {}, " + "firstFile={}, lastFile={}, firstError={}; historical subscription data in these " + "files cannot be replayed"; - public static final String PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_SKIPPED_UNAVAILABLE_SEARCH_INDEXES_B8023B64 = - "ConsensusPrefetchingQueue {}: WAL replay skipped unavailable search indexes [{}, {}), " - + "skippedEntries={}, totalWalGapSkippedEntries={}; the missing WAL data may have been " - + "reclaimed before subscription consumption"; + public static final String PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_FOUND_UNAVAILABLE_SEARCH_INDEXES_E0CBFFFA = + "ConsensusPrefetchingQueue {}: WAL replay found unavailable search indexes [{}, {}) before " + + "searchIndex {}; forcing WAL refresh before retry"; + public static final String MESSAGE_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_CANNOT_RECOVER_SEARCH_INDEXES_70781B22 = + "ConsensusPrefetchingQueue %s: WAL replay cannot recover search indexes [%s, %s), " + + "unavailableEntries=%s, totalWalGapSkippedEntries=%s; subscription delivery is " + + "stalled to prevent silent data loss"; public static final String PIPE_LOG_PIPE_TERMINATE_EVENT_COMMITTED_FOR_HISTORICAL_TRANSFER_CREATIONTIME_9B807B28 = "Pipe {}@{}: terminate event committed for historical transfer. creationTime: {}, " + "shouldMark: {}. {}"; diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 9b124afaffac..b1aca4284795 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2067,9 +2067,12 @@ private DataNodePipeMessages() {} public static final String PIPE_LOG_PROGRESSWALITERATOR_SKIPPED_UNREADABLE_RETAINED_WAL_FILES_FFC8455E = "ProgressWALIterator:跳过了 {} 个无法读取的保留 WAL 文件,directory={},firstFile={}," + "lastFile={},firstError={};这些文件中的历史订阅数据无法重放"; - public static final String PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_SKIPPED_UNAVAILABLE_SEARCH_INDEXES_B8023B64 = - "ConsensusPrefetchingQueue {}:WAL 重放跳过了不可用的 searchIndex 区间 [{}, {})," - + "skippedEntries={},totalWalGapSkippedEntries={};缺失的 WAL 数据可能已在订阅消费前被回收"; + public static final String PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_FOUND_UNAVAILABLE_SEARCH_INDEXES_E0CBFFFA = + "ConsensusPrefetchingQueue {}:WAL 回放发现不可用 searchIndex 区间 [{}, {}),下一个可见 " + + "searchIndex={};正在强制刷新 WAL 后重试"; + public static final String MESSAGE_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_CANNOT_RECOVER_SEARCH_INDEXES_70781B22 = + "ConsensusPrefetchingQueue %s:WAL 回放无法恢复 searchIndex 区间 [%s, %s),不可用条目数=%s," + + "walGapSkippedEntries 总数=%s;为防止静默丢数,订阅投递已停滞"; public static final String PIPE_LOG_PIPE_TERMINATE_EVENT_COMMITTED_FOR_HISTORICAL_TRANSFER_CREATIONTIME_9B807B28 = "Pipe {}@{}:历史传输的终止事件已提交。creationTime:{},shouldMark:{}。{}"; public static final String PIPE_LOG_PIPE_HISTORICAL_SOURCE_HAS_SUPPLIED_ALL_EVENTS_EMITTING_8B58DE19 = diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java index f68d37cc6a24..85807557d979 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java @@ -285,6 +285,12 @@ public class ConsensusPrefetchingQueue { private volatile long lastWalGapWaitLogTimeMs = 0L; + /** Search index whose WAL visibility gap is being retried after forcing a WAL refresh. */ + private volatile long walGapRetryExpectedSearchIndex = Long.MIN_VALUE; + + /** Critical replay failure returned to consumers instead of silently advancing past data. */ + private volatile String walReplayFailureMessage; + /** Fallback committed region progress from local persisted state. */ private final RegionProgress fallbackCommittedRegionProgress; @@ -755,6 +761,9 @@ public SubscriptionEvent poll(final String consumerId, final RegionProgress regi if (pendingSeekRequest != null) { return null; } + if (Objects.nonNull(walReplayFailureMessage)) { + return generateCriticalErrorResponse(walReplayFailureMessage); + } final SubscriptionEvent event = pollInternal(consumerId); if (Objects.nonNull(event) && prefetchingQueue.size() < MAX_PREFETCHING_QUEUE_SIZE) { requestPrefetch(); @@ -821,6 +830,7 @@ private boolean initPrefetchUnderInitializationLock(final RegionProgress regionP // readers without WAL support. this.subscriptionWALIterator = createSubscriptionWALIterator(resolvedStart.getStartSearchIndex()); + resetWalReplayFailureState(); this.prefetchInitialized = true; this.observedSeekGeneration = seekGeneration.get(); discardBatch(this.lingerBatch); @@ -1411,6 +1421,11 @@ public PrefetchRoundResult drivePrefetchOnce() { applyPendingSubscriptionWalReset(observedSeekGeneration); recycleInFlightEvents(); + if (Objects.nonNull(walReplayFailureMessage)) { + blockRealtimeAdmission(); + return PrefetchRoundResult.dormant(); + } + if (!isActive) { blockRealtimeAdmission(); return computeIdleRoundResult(); @@ -1490,6 +1505,11 @@ public PrefetchRoundResult drivePrefetchOnce() { if (batch.isEmpty() && lingerBatch.isEmpty()) { final MaterializationResult walResult = tryCatchUpFromWAL(observedSeekGeneration); + if (walResult == MaterializationResult.WAL_GAP) { + return Objects.nonNull(walReplayFailureMessage) + ? PrefetchRoundResult.dormant() + : PrefetchRoundResult.rescheduleAfter(WAL_GAP_RETRY_SLEEP_MS); + } if (walResult == MaterializationResult.MEMORY_BLOCKED) { blockRealtimeAdmission(); return PrefetchRoundResult.rescheduleAfter(MEMORY_RETRY_SLEEP_MS); @@ -1898,7 +1918,9 @@ private MaterializationResult tryCatchUpFromWAL(final long expectedSeekGeneratio pumpFromSubscriptionWAL( batchState, expectedSeekGeneration, maxWalEntries, maxTablets, maxBatchBytes); if (materializationResult != MaterializationResult.SUCCESS) { - if (materializationResult == MaterializationResult.MEMORY_BLOCKED && !batchState.isEmpty()) { + if ((materializationResult == MaterializationResult.MEMORY_BLOCKED + || materializationResult == MaterializationResult.WAL_GAP) + && !batchState.isEmpty()) { if (!flushBatch(batchState, expectedSeekGeneration)) { discardBatch(batchState); return MaterializationResult.STALE; @@ -1939,6 +1961,10 @@ private MaterializationResult pumpFromSubscriptionWAL( if (isBeforeLocalCursor(walEntry)) { continue; } + final MaterializationResult continuityResult = validateWalReplayContinuity(walEntry); + if (continuityResult != MaterializationResult.SUCCESS) { + return continuityResult; + } if (shouldSkipForRecoveryProgress(walEntry)) { advanceWalReplayCursorIfPresent(walEntry); continue; @@ -1984,22 +2010,67 @@ private void advanceWalReplayCursorIfPresent(final IndexedConsensusRequest reque if (!hasLocalSearchIndex(request)) { return; } + nextExpectedSearchIndex.set(request.getSearchIndex() + 1); + } + + private MaterializationResult validateWalReplayContinuity(final IndexedConsensusRequest request) { + if (!hasLocalSearchIndex(request)) { + return MaterializationResult.SUCCESS; + } final long actualSearchIndex = request.getSearchIndex(); final long expectedSearchIndex = nextExpectedSearchIndex.get(); - if (actualSearchIndex > expectedSearchIndex) { - final long skippedEntries = actualSearchIndex - expectedSearchIndex; - final long totalSkippedEntries = walGapSkippedEntries.addAndGet(skippedEntries); + if (actualSearchIndex <= expectedSearchIndex) { + if (actualSearchIndex == expectedSearchIndex + && walGapRetryExpectedSearchIndex == expectedSearchIndex) { + walGapRetryExpectedSearchIndex = Long.MIN_VALUE; + pendingWalGapRetryRequested = false; + } + return MaterializationResult.SUCCESS; + } + + if (walGapRetryExpectedSearchIndex != expectedSearchIndex) { + walGapRetryExpectedSearchIndex = expectedSearchIndex; LOGGER.warn( DataNodePipeMessages - .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_SKIPPED_UNAVAILABLE_SEARCH_INDEXES_B8023B64, + .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_FOUND_UNAVAILABLE_SEARCH_INDEXES_E0CBFFFA, this, expectedSearchIndex, actualSearchIndex, - skippedEntries, - totalSkippedEntries); + actualSearchIndex); + if (consensusReqReader instanceof WALNode) { + ((WALNode) consensusReqReader).rollWALFile(); + } + resetSubscriptionWALPosition(expectedSearchIndex); + onWalGapRetryScheduled(); + pendingWalGapRetryRequested = true; + return MaterializationResult.WAL_GAP; } - nextExpectedSearchIndex.set(actualSearchIndex + 1); + + if (Objects.isNull(walReplayFailureMessage)) { + final long unavailableEntries = actualSearchIndex - expectedSearchIndex; + final long totalUnavailableEntries = walGapSkippedEntries.addAndGet(unavailableEntries); + walReplayFailureMessage = + String.format( + DataNodePipeMessages + .MESSAGE_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_CANNOT_RECOVER_SEARCH_INDEXES_70781B22, + this, + expectedSearchIndex, + actualSearchIndex, + unavailableEntries, + totalUnavailableEntries); + LOGGER.error(walReplayFailureMessage); + blockRealtimeAdmission(); + } + return MaterializationResult.WAL_GAP; + } + + private void resetWalReplayFailureState() { + pendingWalGapRetryRequested = false; + walGapWaitStartTimeMs = 0L; + lastWalGapWaitLogTimeMs = 0L; + walGapRetryExpectedSearchIndex = Long.MIN_VALUE; + walReplayFailureMessage = null; } private void ensureSubscriptionWalReadable() { @@ -2981,9 +3052,7 @@ public void cleanUp() { reconcileRetainedTabletMemoryAfterCleanup(); memoryBlockedEntryBytes = -1L; resetBatchWriterProgress(); - pendingWalGapRetryRequested = false; - walGapWaitStartTimeMs = 0L; - lastWalGapWaitLogTimeMs = 0L; + resetWalReplayFailureState(); pendingSubscriptionWalResetSearchIndex = Long.MIN_VALUE; pendingSubscriptionWalResetGeneration = Long.MIN_VALUE; closeSubscriptionWALIterator(); @@ -3277,9 +3346,7 @@ private void applySeekResetUnderWriteLock(final PendingSeekRequest request) { reconcileRetainedTabletMemoryAfterCleanup(); resetBatchWriterProgress(); observedSeekGeneration = seekGeneration.get(); - pendingWalGapRetryRequested = false; - walGapWaitStartTimeMs = 0L; - lastWalGapWaitLogTimeMs = 0L; + resetWalReplayFailureState(); // 5. Reset commit state to the writer progress immediately before the first re-delivered // entry so seek/rebind resumes from the intended frontier. @@ -3797,6 +3864,13 @@ private SubscriptionEvent generateErrorResponse(final String errorMessage) { createNonCommittableContext(IoTDBDescriptor.getInstance().getConfig().getDataNodeId())); } + private SubscriptionEvent generateCriticalErrorResponse(final String errorMessage) { + return new SubscriptionEvent( + SubscriptionPollResponseType.ERROR.getType(), + new ErrorPayload(errorMessage, true), + createNonCommittableContext(IoTDBDescriptor.getInstance().getConfig().getDataNodeId())); + } + private SubscriptionEvent generateOutdatedErrorResponse() { return new SubscriptionEvent( SubscriptionPollResponseType.ERROR.getType(), @@ -3901,9 +3975,7 @@ private void setActiveUnderRuntimeLock(final boolean active) { memoryBlockedEntryBytes = -1L; prefetchInitialized = false; observedSeekGeneration = seekGeneration.get(); - pendingWalGapRetryRequested = false; - walGapWaitStartTimeMs = 0L; - lastWalGapWaitLogTimeMs = 0L; + resetWalReplayFailureState(); pendingSubscriptionWalResetSearchIndex = Long.MIN_VALUE; pendingSubscriptionWalResetGeneration = Long.MIN_VALUE; closeSubscriptionWALIterator(); @@ -3991,6 +4063,9 @@ public String getProgressStatusName() { if (!isActive) { return SubscriptionProgressSnapshot.STATUS_INACTIVE; } + if (Objects.nonNull(walReplayFailureMessage)) { + return SubscriptionProgressSnapshot.STATUS_STALLED; + } if (getLag() <= 0L) { return SubscriptionProgressSnapshot.STATUS_CAUGHT_UP; } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java index b21271596726..168a634695cf 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java @@ -44,7 +44,9 @@ import org.apache.iotdb.db.subscription.event.SubscriptionEvent; import org.apache.iotdb.db.subscription.resource.SubscriptionMemoryManager; import org.apache.iotdb.rpc.subscription.config.TopicConstant; +import org.apache.iotdb.rpc.subscription.payload.poll.ErrorPayload; import org.apache.iotdb.rpc.subscription.payload.poll.RegionProgress; +import org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionPollResponseType; import org.apache.iotdb.rpc.subscription.payload.poll.WriterId; import org.apache.iotdb.rpc.subscription.payload.poll.WriterProgress; @@ -1047,27 +1049,125 @@ ProgressWALReader openProgressWALReader(final File walFile) throws IOException { } @Test - public void testWalReplayCountsOnlyUnavailableSearchIndexes() throws Exception { + public void testWalReplayRetriesGapWithoutSkippingEntries() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); final File systemDir = temporaryFolder.newFolder("wal-replay-gap-counter"); ConsensusPrefetchingQueue queue = null; try { final DataRegionId regionId = new DataRegionId(9); + final WALNode walNode = mock(WALNode.class); + when(walNode.getCurrentSearchIndex()).thenReturn(4L); + when(walNode.getLogDirectory()).thenReturn(systemDir); + final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); + when(serverImpl.getConsensusReqReader()).thenReturn(walNode); + when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); + + final AtomicInteger conversionCount = new AtomicInteger(); + final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class); + when(converter.convert(any())) + .thenAnswer( + ignored -> { + conversionCount.incrementAndGet(); + return Collections.singletonList(createTablet()); + }); + when(converter.getDatabaseName()).thenReturn("db"); + + final Iterator initiallyVisibleEntries = + Arrays.asList(createRequest(1L), createRequest(4L)).iterator(); + final ProgressWALIterator initialIterator = mock(ProgressWALIterator.class); + when(initialIterator.hasNext()).thenAnswer(ignored -> initiallyVisibleEntries.hasNext()); + when(initialIterator.next()).thenAnswer(ignored -> initiallyVisibleEntries.next()); + final Iterator refreshedEntries = + Arrays.asList(createRequest(2L), createRequest(3L), createRequest(4L)).iterator(); + final ProgressWALIterator refreshedIterator = mock(ProgressWALIterator.class); + when(refreshedIterator.hasNext()).thenAnswer(ignored -> refreshedEntries.hasNext()); + when(refreshedIterator.next()).thenAnswer(ignored -> refreshedEntries.next()); + final AtomicInteger iteratorCreationCount = new AtomicInteger(); + + queue = + new ConsensusPrefetchingQueue( + "consumerGroup", + "topic", + TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE, + regionId, + serverImpl, + new SubscriptionWalRetentionPolicy( + "topic", + SubscriptionWalRetentionPolicy.UNBOUNDED, + SubscriptionWalRetentionPolicy.UNBOUNDED), + converter, + newCommitManager(systemDir), + new RegionProgress(Collections.emptyMap()), + 1L, + 1L, + true) { + @Override + protected ProgressWALIterator createSubscriptionWALIterator( + final long startSearchIndex) { + return iteratorCreationCount.getAndIncrement() == 0 + ? initialIterator + : refreshedIterator; + } + }; + queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024)); + + assertNull(queue.poll("consumer")); + queue.drivePrefetchOnce(); + + assertEquals(1L, queue.getWalPathAcceptedEntries()); + assertEquals(1, conversionCount.get()); + assertEquals(0L, queue.getWalGapSkippedEntries()); + assertEquals(2L, queue.getCurrentReadSearchIndex()); + verify(walNode).rollWALFile(); + + queue.drivePrefetchOnce(); + + assertEquals(4L, queue.getWalPathAcceptedEntries()); + assertEquals(4, conversionCount.get()); + assertEquals(0L, queue.getWalGapSkippedEntries()); + assertEquals(5L, queue.getCurrentReadSearchIndex()); + } finally { + if (queue != null) { + queue.close(); + } + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + + @Test + public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final File systemDir = temporaryFolder.newFolder("wal-replay-unrecoverable-gap"); + ConsensusPrefetchingQueue queue = null; + try { + final DataRegionId regionId = new DataRegionId(10); final FakeConsensusReqReader reader = new FakeConsensusReqReader(); reader.currentSearchIndex = 4L; final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); when(serverImpl.getConsensusReqReader()).thenReturn(reader); when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); + final AtomicInteger conversionCount = new AtomicInteger(); final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class); - when(converter.convert(any())).thenReturn(Collections.singletonList(createTablet())); + when(converter.convert(any())) + .thenAnswer( + ignored -> { + conversionCount.incrementAndGet(); + return Collections.singletonList(createTablet()); + }); when(converter.getDatabaseName()).thenReturn("db"); - final Iterator retainedWalEntries = + final Iterator initiallyVisibleEntries = Arrays.asList(createRequest(1L), createRequest(4L)).iterator(); - final ProgressWALIterator walIterator = mock(ProgressWALIterator.class); - when(walIterator.hasNext()).thenAnswer(ignored -> retainedWalEntries.hasNext()); - when(walIterator.next()).thenAnswer(ignored -> retainedWalEntries.next()); + final ProgressWALIterator initialIterator = mock(ProgressWALIterator.class); + when(initialIterator.hasNext()).thenAnswer(ignored -> initiallyVisibleEntries.hasNext()); + when(initialIterator.next()).thenAnswer(ignored -> initiallyVisibleEntries.next()); + final Iterator refreshedEntries = + Collections.singletonList(createRequest(4L)).iterator(); + final ProgressWALIterator refreshedIterator = mock(ProgressWALIterator.class); + when(refreshedIterator.hasNext()).thenAnswer(ignored -> refreshedEntries.hasNext()); + when(refreshedIterator.next()).thenAnswer(ignored -> refreshedEntries.next()); + final AtomicInteger iteratorCreationCount = new AtomicInteger(); queue = new ConsensusPrefetchingQueue( @@ -1089,16 +1189,29 @@ public void testWalReplayCountsOnlyUnavailableSearchIndexes() throws Exception { @Override protected ProgressWALIterator createSubscriptionWALIterator( final long startSearchIndex) { - return walIterator; + return iteratorCreationCount.getAndIncrement() == 0 + ? initialIterator + : refreshedIterator; } }; + queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024)); assertNull(queue.poll("consumer")); queue.drivePrefetchOnce(); + queue.drivePrefetchOnce(); - assertEquals(2L, queue.getWalPathAcceptedEntries()); + assertEquals(1L, queue.getWalPathAcceptedEntries()); + assertEquals(1, conversionCount.get()); assertEquals(2L, queue.getWalGapSkippedEntries()); - assertEquals(5L, queue.getCurrentReadSearchIndex()); + assertEquals(2L, queue.getCurrentReadSearchIndex()); + assertEquals(4L, queue.getProgressStatus()); + + final SubscriptionEvent errorEvent = queue.poll("consumer"); + assertNotNull(errorEvent); + assertEquals( + SubscriptionPollResponseType.ERROR.getType(), + errorEvent.getCurrentResponse().getResponseType()); + assertTrue(((ErrorPayload) errorEvent.getCurrentResponse().getPayload()).isCritical()); } finally { if (queue != null) { queue.close(); From 7309615a88e0be2e02758fe4c629dff9f2ab3756 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 29 Sep 2026 15:31:51 +0800 Subject: [PATCH 2/4] Fix transient WAL visibility gaps during subscription replay --- .../consensus/ConsensusPrefetchingQueue.java | 26 +++++ .../ConsensusPrefetchingQueueTest.java | 99 ++++++++++++++++++- 2 files changed, 124 insertions(+), 1 deletion(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java index 85807557d979..c7cd7a5b3636 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java @@ -1397,6 +1397,7 @@ public SubscriptionEvent pollTablets( private static final long WAL_GAP_RETRY_SLEEP_MS = 10L; private static final long MEMORY_RETRY_SLEEP_MS = 100L; private static final long WAL_GAP_WAIT_LOG_INTERVAL_MS = 5_000L; + private static final long WAL_GAP_MAX_WAIT_MS = 30_000L; private static final long PREFETCH_STATS_LOG_INTERVAL_MS = 5_000L; @@ -2031,6 +2032,7 @@ private MaterializationResult validateWalReplayContinuity(final IndexedConsensus if (walGapRetryExpectedSearchIndex != expectedSearchIndex) { walGapRetryExpectedSearchIndex = expectedSearchIndex; + walGapWaitStartTimeMs = System.currentTimeMillis(); LOGGER.warn( DataNodePipeMessages .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_FOUND_UNAVAILABLE_SEARCH_INDEXES_E0CBFFFA, @@ -2047,6 +2049,30 @@ private MaterializationResult validateWalReplayContinuity(final IndexedConsensus return MaterializationResult.WAL_GAP; } + final long nowMs = System.currentTimeMillis(); + final long currentWalSearchIndex = consensusReqReader.getCurrentSearchIndex(); + if (currentWalSearchIndex >= actualSearchIndex + && nowMs - walGapWaitStartTimeMs < WAL_GAP_MAX_WAIT_MS) { + if (lastWalGapWaitLogTimeMs == 0L + || nowMs - lastWalGapWaitLogTimeMs >= WAL_GAP_WAIT_LOG_INTERVAL_MS) { + LOGGER.info( + DataNodePipeMessages + .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAITING_MS_FOR_WAL_GAP_TO_BECOME_7D91C6C5, + this, + nowMs - walGapWaitStartTimeMs, + expectedSearchIndex, + actualSearchIndex, + expectedSearchIndex, + currentWalSearchIndex, + seekGeneration.get()); + lastWalGapWaitLogTimeMs = nowMs; + } + resetSubscriptionWALPosition(expectedSearchIndex); + onWalGapRetryScheduled(); + pendingWalGapRetryRequested = true; + return MaterializationResult.WAL_GAP; + } + if (Objects.isNull(walReplayFailureMessage)) { final long unavailableEntries = actualSearchIndex - expectedSearchIndex; final long totalUnavailableEntries = walGapSkippedEntries.addAndGet(unavailableEntries); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java index 168a634695cf..371b79bca8ae 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java @@ -1134,6 +1134,103 @@ protected ProgressWALIterator createSubscriptionWALIterator( } } + @Test + public void testWalReplayWaitsForWalVisibilityWhenWriterTailIsAhead() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final File systemDir = temporaryFolder.newFolder("wal-replay-gap-wal-catch-up"); + ConsensusPrefetchingQueue queue = null; + try { + final DataRegionId regionId = new DataRegionId(11); + final WALNode walNode = mock(WALNode.class); + when(walNode.getCurrentSearchIndex()).thenReturn(4L); + when(walNode.getLogDirectory()).thenReturn(systemDir); + final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); + when(serverImpl.getConsensusReqReader()).thenReturn(walNode); + when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); + + final AtomicInteger conversionCount = new AtomicInteger(); + final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class); + when(converter.convert(any())) + .thenAnswer( + ignored -> { + conversionCount.incrementAndGet(); + return Collections.singletonList(createTablet()); + }); + when(converter.getDatabaseName()).thenReturn("db"); + + final ProgressWALIterator initialIterator = mock(ProgressWALIterator.class); + final Iterator initiallyVisibleEntries = + Arrays.asList(createRequest(1L), createRequest(4L)).iterator(); + when(initialIterator.hasNext()).thenAnswer(ignored -> initiallyVisibleEntries.hasNext()); + when(initialIterator.next()).thenAnswer(ignored -> initiallyVisibleEntries.next()); + + final ProgressWALIterator laggingIterator = mock(ProgressWALIterator.class); + final Iterator laggingEntries = + Collections.singletonList(createRequest(4L)).iterator(); + when(laggingIterator.hasNext()).thenAnswer(ignored -> laggingEntries.hasNext()); + when(laggingIterator.next()).thenAnswer(ignored -> laggingEntries.next()); + + final ProgressWALIterator recoveredIterator = mock(ProgressWALIterator.class); + final Iterator recoveredEntries = + Arrays.asList(createRequest(2L), createRequest(3L), createRequest(4L)).iterator(); + when(recoveredIterator.hasNext()).thenAnswer(ignored -> recoveredEntries.hasNext()); + when(recoveredIterator.next()).thenAnswer(ignored -> recoveredEntries.next()); + + final AtomicInteger iteratorCreationCount = new AtomicInteger(); + queue = + new ConsensusPrefetchingQueue( + "consumerGroup", + "topic", + TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE, + regionId, + serverImpl, + new SubscriptionWalRetentionPolicy( + "topic", + SubscriptionWalRetentionPolicy.UNBOUNDED, + SubscriptionWalRetentionPolicy.UNBOUNDED), + converter, + newCommitManager(systemDir), + new RegionProgress(Collections.emptyMap()), + 1L, + 1L, + true) { + @Override + protected ProgressWALIterator createSubscriptionWALIterator( + final long startSearchIndex) { + switch (iteratorCreationCount.getAndIncrement()) { + case 0: + return initialIterator; + case 1: + return laggingIterator; + default: + return recoveredIterator; + } + } + }; + queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024)); + + assertNull(queue.poll("consumer")); + queue.drivePrefetchOnce(); + assertEquals(1L, queue.getWalPathAcceptedEntries()); + assertEquals(0L, queue.getWalGapSkippedEntries()); + + queue.drivePrefetchOnce(); + assertEquals(1L, queue.getWalPathAcceptedEntries()); + assertEquals(0L, queue.getWalGapSkippedEntries()); + + queue.drivePrefetchOnce(); + assertEquals(4L, queue.getWalPathAcceptedEntries()); + assertEquals(4, conversionCount.get()); + assertEquals(0L, queue.getWalGapSkippedEntries()); + verify(walNode).rollWALFile(); + } finally { + if (queue != null) { + queue.close(); + } + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -1142,7 +1239,7 @@ public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Excep try { final DataRegionId regionId = new DataRegionId(10); final FakeConsensusReqReader reader = new FakeConsensusReqReader(); - reader.currentSearchIndex = 4L; + reader.currentSearchIndex = 3L; final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); when(serverImpl.getConsensusReqReader()).thenReturn(reader); when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); From 9fea927c13222d8d1762be29ab04ef87f77a3e1c Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 29 Sep 2026 15:55:47 +0800 Subject: [PATCH 3/4] Fix subscription WAL retention at replay cursor --- .../consensus/iot/IoTConsensusServerImpl.java | 7 +- .../SubscriptionQueueRegistry.java | 21 +- .../SubscriptionWalRetentionCalculator.java | 10 +- .../consensus/ConsensusPrefetchingQueue.java | 56 ++--- .../ConsensusPrefetchingQueueTest.java | 196 +++++++++--------- 5 files changed, 145 insertions(+), 145 deletions(-) diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java index f4448e81f34b..ba4ad87b1fdf 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java @@ -1266,9 +1266,8 @@ public ConsensusReqReader getConsensusReqReader() { public void registerSubscriptionQueue( final BlockingQueue queue, final SubscriptionWalRetentionPolicy retentionPolicy, - final LongSupplier committedRetainedMinVersionIdSupplier) { - subscriptionQueueRegistry.register( - queue, retentionPolicy, committedRetainedMinVersionIdSupplier); + final LongSupplier retainedMinVersionIdSupplier) { + subscriptionQueueRegistry.register(queue, retentionPolicy, retainedMinVersionIdSupplier); // Immediately re-evaluate the safe delete index with new subscription awareness checkAndUpdateSafeDeletedSearchIndex(); logger.info( @@ -1440,7 +1439,7 @@ public void checkAndUpdateSafeDeletedSearchIndex() { final SubscriptionRetentionBound subscriptionRetentionBound = subscriptionWalRetentionCalculator.calculate( subscriptionQueueRegistry.getRetentionPolicies(), - subscriptionQueueRegistry.getCommittedRetainedMinVersionIds()); + subscriptionQueueRegistry.getRetainedMinVersionIds()); consensusReqReader.setSafelyDeletedSearchIndex( Math.min(replicationIndex, subscriptionRetentionBound.getSafelyDeletedSearchIndex())); diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java index 8d7c27f8cf2a..d63e42f12c98 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java @@ -67,10 +67,9 @@ public synchronized void register( public synchronized void register( final BlockingQueue queue, final SubscriptionWalRetentionPolicy retentionPolicy, - final LongSupplier committedRetainedMinVersionIdSupplier) { + final LongSupplier retainedMinVersionIdSupplier) { queues.put( - queue, - new SubscriptionQueueRegistration(retentionPolicy, committedRetainedMinVersionIdSupplier)); + queue, new SubscriptionQueueRegistration(retentionPolicy, retainedMinVersionIdSupplier)); } // Shares the monitor with offer() so unregister() is a real stop-receiving barrier. @@ -94,19 +93,19 @@ public synchronized Collection getRetentionPolic return retentionPolicies; } - public Collection getCommittedRetainedMinVersionIds() { + public Collection getRetainedMinVersionIds() { final Collection suppliers = new ArrayList<>(); synchronized (this) { for (final SubscriptionQueueRegistration registration : queues.values()) { - suppliers.add(registration.committedRetainedMinVersionIdSupplier); + suppliers.add(registration.retainedMinVersionIdSupplier); } } - final Collection committedRetainedMinVersionIds = new ArrayList<>(); + final Collection retainedMinVersionIds = new ArrayList<>(); for (final LongSupplier supplier : suppliers) { - committedRetainedMinVersionIds.add(supplier.getAsLong()); + retainedMinVersionIds.add(supplier.getAsLong()); } - return committedRetainedMinVersionIds; + return retainedMinVersionIds; } public synchronized boolean offer(final IndexedConsensusRequest indexedConsensusRequest) { @@ -172,13 +171,13 @@ public synchronized boolean offer(final IndexedConsensusRequest indexedConsensus private static final class SubscriptionQueueRegistration { private final SubscriptionWalRetentionPolicy retentionPolicy; - private final LongSupplier committedRetainedMinVersionIdSupplier; + private final LongSupplier retainedMinVersionIdSupplier; private SubscriptionQueueRegistration( final SubscriptionWalRetentionPolicy retentionPolicy, - final LongSupplier committedRetainedMinVersionIdSupplier) { + final LongSupplier retainedMinVersionIdSupplier) { this.retentionPolicy = retentionPolicy; - this.committedRetainedMinVersionIdSupplier = committedRetainedMinVersionIdSupplier; + this.retainedMinVersionIdSupplier = retainedMinVersionIdSupplier; } } } diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java index 44c0ec47f15e..5845e2686759 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java @@ -85,7 +85,7 @@ public SubscriptionWalRetentionCalculator(final ConsensusReqReader consensusReqR public SubscriptionRetentionBound calculate( final Collection retentionPolicies, - final Collection committedRetainedMinVersionIds) { + final Collection queueRetainedMinVersionIds) { SubscriptionRetentionBound mergedBound = SubscriptionRetentionBound.noConstraint(); for (final SubscriptionWalRetentionPolicy policy : retentionPolicies) { // For each topic, data can be deleted once either its size retention or its time retention @@ -96,14 +96,14 @@ public SubscriptionRetentionBound calculate( .mergeDeleteEither(buildTimeRetentionBound(policy.getRetentionMs())); mergedBound = mergedBound.mergeDeleteOnlyIfBoth(perQueueBound); } - for (final long committedRetainedMinVersionId : committedRetainedMinVersionIds) { - // Topic retention is a historical replay window, while committed progress protects data - // that a consumer group has not acknowledged yet. A WAL file can be reclaimed only when + for (final long queueRetainedMinVersionId : queueRetainedMinVersionIds) { + // Topic retention is a historical replay window, while each queue also protects WAL needed + // by its committed progress and local replay cursor. A WAL file can be reclaimed only when // both constraints allow it, so keep the more conservative (smaller) file-version bound. mergedBound = mergedBound.mergeDeleteOnlyIfBoth( SubscriptionRetentionBound.of( - Long.MAX_VALUE, Math.max(0L, committedRetainedMinVersionId))); + Long.MAX_VALUE, Math.max(0L, queueRetainedMinVersionId))); } return mergedBound; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java index c7cd7a5b3636..95c890a4d47a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java @@ -573,7 +573,7 @@ public ConsensusPrefetchingQueue( this::requestPrefetchForRealtimeEntry, this::canAcceptRealtimeEntry); serverImpl.registerSubscriptionQueue( - pendingEntries, retentionPolicy, this::getCommittedRetainedMinVersionId); + pendingEntries, retentionPolicy, this::getRequiredRetainedMinVersionId); LOGGER.info( DataNodePipeMessages @@ -1397,7 +1397,6 @@ public SubscriptionEvent pollTablets( private static final long WAL_GAP_RETRY_SLEEP_MS = 10L; private static final long MEMORY_RETRY_SLEEP_MS = 100L; private static final long WAL_GAP_WAIT_LOG_INTERVAL_MS = 5_000L; - private static final long WAL_GAP_MAX_WAIT_MS = 30_000L; private static final long PREFETCH_STATS_LOG_INTERVAL_MS = 5_000L; @@ -2032,7 +2031,6 @@ private MaterializationResult validateWalReplayContinuity(final IndexedConsensus if (walGapRetryExpectedSearchIndex != expectedSearchIndex) { walGapRetryExpectedSearchIndex = expectedSearchIndex; - walGapWaitStartTimeMs = System.currentTimeMillis(); LOGGER.warn( DataNodePipeMessages .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_FOUND_UNAVAILABLE_SEARCH_INDEXES_E0CBFFFA, @@ -2049,30 +2047,6 @@ private MaterializationResult validateWalReplayContinuity(final IndexedConsensus return MaterializationResult.WAL_GAP; } - final long nowMs = System.currentTimeMillis(); - final long currentWalSearchIndex = consensusReqReader.getCurrentSearchIndex(); - if (currentWalSearchIndex >= actualSearchIndex - && nowMs - walGapWaitStartTimeMs < WAL_GAP_MAX_WAIT_MS) { - if (lastWalGapWaitLogTimeMs == 0L - || nowMs - lastWalGapWaitLogTimeMs >= WAL_GAP_WAIT_LOG_INTERVAL_MS) { - LOGGER.info( - DataNodePipeMessages - .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAITING_MS_FOR_WAL_GAP_TO_BECOME_7D91C6C5, - this, - nowMs - walGapWaitStartTimeMs, - expectedSearchIndex, - actualSearchIndex, - expectedSearchIndex, - currentWalSearchIndex, - seekGeneration.get()); - lastWalGapWaitLogTimeMs = nowMs; - } - resetSubscriptionWALPosition(expectedSearchIndex); - onWalGapRetryScheduled(); - pendingWalGapRetryRequested = true; - return MaterializationResult.WAL_GAP; - } - if (Objects.isNull(walReplayFailureMessage)) { final long unavailableEntries = actualSearchIndex - expectedSearchIndex; final long totalUnavailableEntries = walGapSkippedEntries.addAndGet(unavailableEntries); @@ -3416,11 +3390,39 @@ private boolean commitWithoutOutstandingAndRefreshWalRetention( return committed; } + private long getRequiredRetainedMinVersionId() { + // Committed progress can advance on another consumer or replica while this queue is still + // replaying an older local search index. Keep both boundaries so WAL deletion can never pass + // the earliest entry this queue has not inspected yet. + return Math.min(getCommittedRetainedMinVersionId(), getReplayRetainedMinVersionId()); + } + private long getCommittedRetainedMinVersionId() { refreshCommittedWalRetentionBound(); return committedRetainedMinVersionId; } + private long getReplayRetainedMinVersionId() { + if (!(consensusReqReader instanceof WALNode)) { + return 0L; + } + final WALNode walNode = (WALNode) consensusReqReader; + return findReplayRetainedMinVersionId( + WALFileUtils.listAllWALFiles(walNode.getLogDirectory()), nextExpectedSearchIndex.get()); + } + + static long findReplayRetainedMinVersionId( + final File[] walFiles, final long nextExpectedSearchIndex) { + if (Objects.isNull(walFiles) || walFiles.length == 0) { + return 0L; + } + + WALFileUtils.ascSortByVersionId(walFiles); + final int replayFileIndex = + Math.max(0, WALFileUtils.binarySearchFileBySearchIndex(walFiles, nextExpectedSearchIndex)); + return WALFileUtils.parseVersionId(walFiles[replayFileIndex].getName()); + } + private void refreshCommittedWalRetentionBoundAndNotify() { if (refreshCommittedWalRetentionBound()) { serverImpl.checkAndUpdateSafeDeletedSearchIndex(); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java index 371b79bca8ae..f86eae20abe9 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java @@ -80,6 +80,7 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.function.BooleanSupplier; +import java.util.function.LongSupplier; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -223,6 +224,88 @@ public void testWalFileCommitRequirementRejectsUnsupportedWriterMetadata() { assertFalse(requirement.isCoveredBy(new RegionProgress(Collections.emptyMap()))); } + @Test + public void testReplayCursorBoundsWalRetentionWhenCommittedProgressIsAhead() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final File systemDir = temporaryFolder.newFolder("replay-cursor-retention-system"); + final File walDirectory = temporaryFolder.newFolder("replay-cursor-retention-wal"); + ConsensusPrefetchingQueue queue = null; + try { + final File firstWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(0L, 0L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + final File unreadWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(1L, 100L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + final File liveWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(2L, 200L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + writeWalMetadata(firstWal, 1L, 1L, 1001L, 7); + writeWalMetadata(unreadWal, 101L, 101L, 1101L, 7); + try (WALWriter ignored = new WALWriter(liveWal, WALFileVersion.V3)) { + // Keep an empty live successor so both data files are eligible for deletion. + } + + final DataRegionId regionId = new DataRegionId(1); + final WALNode walNode = mock(WALNode.class); + when(walNode.getLogDirectory()).thenReturn(walDirectory); + when(walNode.getCurrentWALFileVersion()).thenReturn(2L); + + final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); + when(serverImpl.getConsensusReqReader()).thenReturn(walNode); + when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); + final AtomicReference retentionSupplier = new AtomicReference<>(); + doAnswer( + invocation -> { + retentionSupplier.set(invocation.getArgument(2)); + return null; + }) + .when(serverImpl) + .registerSubscriptionQueue(any(), any(), any()); + + final ConsensusSubscriptionCommitManager commitManager = + mock(ConsensusSubscriptionCommitManager.class); + when(commitManager.getCommittedRegionProgress("consumerGroup", "topic", regionId)) + .thenReturn( + new RegionProgress( + Collections.singletonMap( + new WriterId(regionId.toString(), 7), new WriterProgress(1101L, 101L)))); + + queue = + new ConsensusPrefetchingQueue( + "consumerGroup", + "topic", + TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE, + regionId, + serverImpl, + new SubscriptionWalRetentionPolicy( + "topic", + SubscriptionWalRetentionPolicy.UNBOUNDED, + SubscriptionWalRetentionPolicy.UNBOUNDED), + mock(ConsensusLogToTabletConverter.class), + commitManager, + new RegionProgress(Collections.emptyMap()), + 101L, + 1L, + true); + + assertNotNull(retentionSupplier.get()); + assertEquals(1L, retentionSupplier.get().getAsLong()); + final Field committedBound = + ConsensusPrefetchingQueue.class.getDeclaredField("committedRetainedMinVersionId"); + committedBound.setAccessible(true); + assertEquals(2L, committedBound.getLong(queue)); + } finally { + if (queue != null) { + queue.close(); + } + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testReplayStartPreservesUncoveredFollowerEntries() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -1134,103 +1217,6 @@ protected ProgressWALIterator createSubscriptionWALIterator( } } - @Test - public void testWalReplayWaitsForWalVisibilityWhenWriterTailIsAhead() throws Exception { - final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); - final File systemDir = temporaryFolder.newFolder("wal-replay-gap-wal-catch-up"); - ConsensusPrefetchingQueue queue = null; - try { - final DataRegionId regionId = new DataRegionId(11); - final WALNode walNode = mock(WALNode.class); - when(walNode.getCurrentSearchIndex()).thenReturn(4L); - when(walNode.getLogDirectory()).thenReturn(systemDir); - final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); - when(serverImpl.getConsensusReqReader()).thenReturn(walNode); - when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); - - final AtomicInteger conversionCount = new AtomicInteger(); - final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class); - when(converter.convert(any())) - .thenAnswer( - ignored -> { - conversionCount.incrementAndGet(); - return Collections.singletonList(createTablet()); - }); - when(converter.getDatabaseName()).thenReturn("db"); - - final ProgressWALIterator initialIterator = mock(ProgressWALIterator.class); - final Iterator initiallyVisibleEntries = - Arrays.asList(createRequest(1L), createRequest(4L)).iterator(); - when(initialIterator.hasNext()).thenAnswer(ignored -> initiallyVisibleEntries.hasNext()); - when(initialIterator.next()).thenAnswer(ignored -> initiallyVisibleEntries.next()); - - final ProgressWALIterator laggingIterator = mock(ProgressWALIterator.class); - final Iterator laggingEntries = - Collections.singletonList(createRequest(4L)).iterator(); - when(laggingIterator.hasNext()).thenAnswer(ignored -> laggingEntries.hasNext()); - when(laggingIterator.next()).thenAnswer(ignored -> laggingEntries.next()); - - final ProgressWALIterator recoveredIterator = mock(ProgressWALIterator.class); - final Iterator recoveredEntries = - Arrays.asList(createRequest(2L), createRequest(3L), createRequest(4L)).iterator(); - when(recoveredIterator.hasNext()).thenAnswer(ignored -> recoveredEntries.hasNext()); - when(recoveredIterator.next()).thenAnswer(ignored -> recoveredEntries.next()); - - final AtomicInteger iteratorCreationCount = new AtomicInteger(); - queue = - new ConsensusPrefetchingQueue( - "consumerGroup", - "topic", - TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE, - regionId, - serverImpl, - new SubscriptionWalRetentionPolicy( - "topic", - SubscriptionWalRetentionPolicy.UNBOUNDED, - SubscriptionWalRetentionPolicy.UNBOUNDED), - converter, - newCommitManager(systemDir), - new RegionProgress(Collections.emptyMap()), - 1L, - 1L, - true) { - @Override - protected ProgressWALIterator createSubscriptionWALIterator( - final long startSearchIndex) { - switch (iteratorCreationCount.getAndIncrement()) { - case 0: - return initialIterator; - case 1: - return laggingIterator; - default: - return recoveredIterator; - } - } - }; - queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024)); - - assertNull(queue.poll("consumer")); - queue.drivePrefetchOnce(); - assertEquals(1L, queue.getWalPathAcceptedEntries()); - assertEquals(0L, queue.getWalGapSkippedEntries()); - - queue.drivePrefetchOnce(); - assertEquals(1L, queue.getWalPathAcceptedEntries()); - assertEquals(0L, queue.getWalGapSkippedEntries()); - - queue.drivePrefetchOnce(); - assertEquals(4L, queue.getWalPathAcceptedEntries()); - assertEquals(4, conversionCount.get()); - assertEquals(0L, queue.getWalGapSkippedEntries()); - verify(walNode).rollWALFile(); - } finally { - if (queue != null) { - queue.close(); - } - IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); - } - } - @Test public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -1239,7 +1225,7 @@ public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Excep try { final DataRegionId regionId = new DataRegionId(10); final FakeConsensusReqReader reader = new FakeConsensusReqReader(); - reader.currentSearchIndex = 3L; + reader.currentSearchIndex = 4L; final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); when(serverImpl.getConsensusReqReader()).thenReturn(reader); when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); @@ -2530,6 +2516,20 @@ private static IndexedConsensusRequest createRequest(final long searchIndex) { .setNodeId(7); } + private static void writeWalMetadata( + final File walFile, + final long searchIndex, + final long localSeq, + final long physicalTime, + final int writerNodeId) + throws IOException { + final WALMetaData metadata = new WALMetaData(); + metadata.add(1, searchIndex, 1L, physicalTime, writerNodeId, localSeq); + try (WALWriter writer = new WALWriter(walFile, WALFileVersion.V3)) { + writer.write(ByteBuffer.wrap(new byte[] {0}), metadata); + } + } + private static IndexedConsensusRequest createRequest( final long searchIndex, final long localSeq, From 467c7a58685a4cf66f4f58ca4293ace2196c3463 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 30 Sep 2026 19:00:59 +0800 Subject: [PATCH 4/4] Fix duplicate WAL metadata test helper --- .../consensus/ConsensusPrefetchingQueueTest.java | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java index a3668b72e36f..a854b8ad6dc7 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java @@ -2946,20 +2946,6 @@ private static IndexedConsensusRequest createRequest(final long searchIndex) { .setNodeId(7); } - private static void writeWalMetadata( - final File walFile, - final long searchIndex, - final long localSeq, - final long physicalTime, - final int writerNodeId) - throws IOException { - final WALMetaData metadata = new WALMetaData(); - metadata.add(1, searchIndex, 1L, physicalTime, writerNodeId, localSeq); - try (WALWriter writer = new WALWriter(walFile, WALFileVersion.V3)) { - writer.write(ByteBuffer.wrap(new byte[] {0}), metadata); - } - } - private static IndexedConsensusRequest createRequest( final long searchIndex, final long localSeq,