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 1c0362b627b3..95583ded3a92 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 @@ -2216,10 +2216,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 49096659ade1..ccb9cd42da20 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 @@ -2057,9 +2057,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 0dcddbef7a23..c5641db1f182 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 @@ -297,6 +297,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; @@ -771,6 +777,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(); @@ -837,6 +846,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); @@ -1427,6 +1437,11 @@ public PrefetchRoundResult drivePrefetchOnce() { applyPendingSubscriptionWalReset(observedSeekGeneration); recycleInFlightEvents(); + if (Objects.nonNull(walReplayFailureMessage)) { + blockRealtimeAdmission(); + return PrefetchRoundResult.dormant(); + } + if (!isActive) { blockRealtimeAdmission(); return computeIdleRoundResult(); @@ -1511,6 +1526,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); @@ -1926,7 +1946,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; @@ -1969,6 +1991,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; @@ -2031,22 +2057,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; + } + + 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(); } - nextExpectedSearchIndex.set(actualSearchIndex + 1); + return MaterializationResult.WAL_GAP; + } + + private void resetWalReplayFailureState() { + pendingWalGapRetryRequested = false; + walGapWaitStartTimeMs = 0L; + lastWalGapWaitLogTimeMs = 0L; + walGapRetryExpectedSearchIndex = Long.MIN_VALUE; + walReplayFailureMessage = null; } private void ensureSubscriptionWalReadable() { @@ -3049,9 +3120,7 @@ public void cleanUp() { reconcileRetainedTabletMemoryAfterCleanup(); memoryBlockedEntryBytes = -1L; resetBatchWriterProgress(); - pendingWalGapRetryRequested = false; - walGapWaitStartTimeMs = 0L; - lastWalGapWaitLogTimeMs = 0L; + resetWalReplayFailureState(); pendingSubscriptionWalResetSearchIndex = Long.MIN_VALUE; pendingSubscriptionWalResetGeneration = Long.MIN_VALUE; closeSubscriptionWALIterator(); @@ -3345,9 +3414,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. @@ -3448,6 +3515,7 @@ static long findReplayRetainedMinVersionId( return 0L; } + WALFileUtils.ascSortByVersionId(walFiles); final int replayFileIndex = Math.max(0, WALFileUtils.binarySearchFileBySearchIndex(walFiles, nextExpectedSearchIndex)); return WALFileUtils.parseVersionId(walFiles[replayFileIndex].getName()); @@ -3960,6 +4028,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(), @@ -4064,9 +4139,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(); @@ -4154,6 +4227,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 6efb3298bf90..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 @@ -1382,6 +1382,92 @@ ProgressWALReader openProgressWALReader(final File walFile) throws IOException { } } + @Test + 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 testPendingGapReplayHonorsPerRoundWalEntryLimit() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -1523,6 +1609,92 @@ public void testWalReplayCountsOnlyUnavailableSearchIndexes() throws Exception { } } + @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())) + .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 = + 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( + "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(); + queue.drivePrefetchOnce(); + + assertEquals(1L, queue.getWalPathAcceptedEntries()); + assertEquals(1, conversionCount.get()); + assertEquals(2L, queue.getWalGapSkippedEntries()); + 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(); + } + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testActivationRetriesUntilConfigNodeProgressIsExplicitlyAvailable() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();