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 95583ded3a92..1c0362b627b3 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,13 +2216,10 @@ 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_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_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_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 ccb9cd42da20..49096659ade1 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,12 +2057,9 @@ 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_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_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_SKIPPED_UNAVAILABLE_SEARCH_INDEXES_B8023B64 = + "ConsensusPrefetchingQueue {}:WAL 重放跳过了不可用的 searchIndex 区间 [{}, {})," + + "skippedEntries={},totalWalGapSkippedEntries={};缺失的 WAL 数据可能已在订阅消费前被回收"; 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 fa01673ccc03..b692ec94ed04 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,12 +297,6 @@ 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; @@ -777,9 +771,6 @@ 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(); @@ -846,7 +837,6 @@ 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); @@ -1437,11 +1427,6 @@ public PrefetchRoundResult drivePrefetchOnce() { applyPendingSubscriptionWalReset(observedSeekGeneration); recycleInFlightEvents(); - if (Objects.nonNull(walReplayFailureMessage)) { - blockRealtimeAdmission(); - return PrefetchRoundResult.dormant(); - } - if (!isActive) { blockRealtimeAdmission(); return computeIdleRoundResult(); @@ -1530,11 +1515,6 @@ 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); @@ -1950,9 +1930,7 @@ private MaterializationResult tryCatchUpFromWAL(final long expectedSeekGeneratio pumpFromSubscriptionWAL( batchState, expectedSeekGeneration, maxWalEntries, maxTablets, maxBatchBytes); if (materializationResult != MaterializationResult.SUCCESS) { - if ((materializationResult == MaterializationResult.MEMORY_BLOCKED - || materializationResult == MaterializationResult.WAL_GAP) - && !batchState.isEmpty()) { + if (materializationResult == MaterializationResult.MEMORY_BLOCKED && !batchState.isEmpty()) { if (!flushBatch(batchState, expectedSeekGeneration)) { discardBatch(batchState); return MaterializationResult.STALE; @@ -1995,10 +1973,6 @@ private MaterializationResult pumpFromSubscriptionWAL( if (isBeforeLocalCursor(walEntry)) { continue; } - final MaterializationResult continuityResult = validateWalReplayContinuity(walEntry); - if (continuityResult != MaterializationResult.SUCCESS) { - return continuityResult; - } if (shouldSkipForRecoveryProgress(walEntry)) { advanceWalReplayCursorIfPresent(walEntry); continue; @@ -2061,67 +2035,22 @@ 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) { - if (actualSearchIndex == expectedSearchIndex - && walGapRetryExpectedSearchIndex == expectedSearchIndex) { - walGapRetryExpectedSearchIndex = Long.MIN_VALUE; - pendingWalGapRetryRequested = false; - } - return MaterializationResult.SUCCESS; - } - - if (walGapRetryExpectedSearchIndex != expectedSearchIndex) { - walGapRetryExpectedSearchIndex = expectedSearchIndex; + if (actualSearchIndex > expectedSearchIndex) { + final long skippedEntries = actualSearchIndex - expectedSearchIndex; + final long totalSkippedEntries = walGapSkippedEntries.addAndGet(skippedEntries); LOGGER.warn( DataNodePipeMessages - .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_FOUND_UNAVAILABLE_SEARCH_INDEXES_E0CBFFFA, + .PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_WAL_REPLAY_SKIPPED_UNAVAILABLE_SEARCH_INDEXES_B8023B64, this, expectedSearchIndex, actualSearchIndex, - 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(); + skippedEntries, + totalSkippedEntries); } - return MaterializationResult.WAL_GAP; - } - - private void resetWalReplayFailureState() { - pendingWalGapRetryRequested = false; - walGapWaitStartTimeMs = 0L; - lastWalGapWaitLogTimeMs = 0L; - walGapRetryExpectedSearchIndex = Long.MIN_VALUE; - walReplayFailureMessage = null; + nextExpectedSearchIndex.set(actualSearchIndex + 1); } private void ensureSubscriptionWalReadable() { @@ -3124,7 +3053,9 @@ public void cleanUp() { reconcileRetainedTabletMemoryAfterCleanup(); memoryBlockedEntryBytes = -1L; resetBatchWriterProgress(); - resetWalReplayFailureState(); + pendingWalGapRetryRequested = false; + walGapWaitStartTimeMs = 0L; + lastWalGapWaitLogTimeMs = 0L; pendingSubscriptionWalResetSearchIndex = Long.MIN_VALUE; pendingSubscriptionWalResetGeneration = Long.MIN_VALUE; closeSubscriptionWALIterator(); @@ -3418,7 +3349,9 @@ private void applySeekResetUnderWriteLock(final PendingSeekRequest request) { reconcileRetainedTabletMemoryAfterCleanup(); resetBatchWriterProgress(); observedSeekGeneration = seekGeneration.get(); - resetWalReplayFailureState(); + pendingWalGapRetryRequested = false; + walGapWaitStartTimeMs = 0L; + lastWalGapWaitLogTimeMs = 0L; // 5. Reset commit state to the writer progress immediately before the first re-delivered // entry so seek/rebind resumes from the intended frontier. @@ -3519,7 +3452,6 @@ static long findReplayRetainedMinVersionId( return 0L; } - WALFileUtils.ascSortByVersionId(walFiles); final int replayFileIndex = Math.max(0, WALFileUtils.binarySearchFileBySearchIndex(walFiles, nextExpectedSearchIndex)); return WALFileUtils.parseVersionId(walFiles[replayFileIndex].getName()); @@ -4032,13 +3964,6 @@ 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(), @@ -4143,7 +4068,9 @@ private void setActiveUnderRuntimeLock(final boolean active) { memoryBlockedEntryBytes = -1L; prefetchInitialized = false; observedSeekGeneration = seekGeneration.get(); - resetWalReplayFailureState(); + pendingWalGapRetryRequested = false; + walGapWaitStartTimeMs = 0L; + lastWalGapWaitLogTimeMs = 0L; pendingSubscriptionWalResetSearchIndex = Long.MIN_VALUE; pendingSubscriptionWalResetGeneration = Long.MIN_VALUE; closeSubscriptionWALIterator(); @@ -4231,9 +4158,6 @@ 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 a23d9a1d16ce..0f321341401e 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,92 +1382,6 @@ 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(); @@ -1697,92 +1611,6 @@ protected ProgressWALIterator createSubscriptionWALIterator( } } - @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();