Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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: {}. {}";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -1437,11 +1427,6 @@ public PrefetchRoundResult drivePrefetchOnce() {
applyPendingSubscriptionWalReset(observedSeekGeneration);
recycleInFlightEvents();

if (Objects.nonNull(walReplayFailureMessage)) {
blockRealtimeAdmission();
return PrefetchRoundResult.dormant();
}

if (!isActive) {
blockRealtimeAdmission();
return computeIdleRoundResult();
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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());
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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;
}
Expand Down
Loading
Loading