Skip to content
Open
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 @@ -24,6 +24,7 @@
import org.apache.iotdb.commons.client.ThriftClient;
import org.apache.iotdb.commons.client.async.AsyncPipeDataTransferServiceClient;
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkCriticalException;
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkNonReportTimeConfigurableException;
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkResourceException;
import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
Expand Down Expand Up @@ -119,6 +120,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
"Failed to retry transferring events in the retry queue. Remaining events: %d (tablet events: %d, tsfile events: %d). Last failure: %s.";

private static final boolean isSplitTSFileBatchModeEnabled = true;
private static final int MAX_FILE_NOT_FOUND_RETRY_TIMES = 5;

private final IoTDBDataRegionSyncSink syncConnector = new IoTDBDataRegionSyncSink();
private final BlockingQueue<Event> retryEventQueue = new LinkedBlockingQueue<>();
Expand All @@ -129,6 +131,10 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
// semantics because the same payload may compare equal.
private final Map<Event, PipeResourceFailureType> retryEvent2ResourceFailureType =
new IdentityHashMap<>();
// File disappearance is permanent for a given TsFile event. Keep its retry count separate from
// transient receiver failures so one missing file cannot retry forever in the async queue.
private final Map<Event, Integer> missingFileRetryTimes = new IdentityHashMap<>();
private volatile PipeRuntimeSinkCriticalException missingFileRetryLimitException;
// Keep only the latest text to avoid retaining the complete exception chain for every event.
private volatile String lastRetryFailureMessage;
// Events removed from the retry queue remain here while their next transfer is being started.
Expand Down Expand Up @@ -510,6 +516,21 @@ private boolean transferWithoutCheck(final TsFileInsertionEvent tsFileInsertionE
throw new FileNotFoundException(pipeTsFileInsertionEvent.getTsFile().getAbsolutePath());
}

final boolean transferMod =
pipeTsFileInsertionEvent.isWithMod()
&& clientManager.supportModsIfIsDataNodeReceiver()
&& pipeTsFileInsertionEvent.getModFile() != null
&& pipeTsFileInsertionEvent.getModFile().exists();
if (pipeTsFileInsertionEvent.isWithMod()
&& clientManager.supportModsIfIsDataNodeReceiver()
&& pipeTsFileInsertionEvent.getModFile() != null
&& !pipeTsFileInsertionEvent.getModFile().exists()) {
LOGGER.warn(
"TsFile {}: modification file {} is missing, transfer the TsFile without modifications.",
pipeTsFileInsertionEvent.getTsFile(),
pipeTsFileInsertionEvent.getModFile());
}

pipeTransferTsFileHandler =
new PipeTransferTsFileHandler(
this,
Expand All @@ -523,8 +544,7 @@ private boolean transferWithoutCheck(final TsFileInsertionEvent tsFileInsertionE
new AtomicBoolean(false),
pipeTsFileInsertionEvent.getTsFile(),
pipeTsFileInsertionEvent.getModFile(),
pipeTsFileInsertionEvent.isWithMod()
&& clientManager.supportModsIfIsDataNodeReceiver(),
transferMod,
pipeTsFileInsertionEvent.getDatabaseName());
trackRetryHandler(
pipeTransferTsFileHandler, Collections.singletonList(pipeTsFileInsertionEvent));
Expand Down Expand Up @@ -661,6 +681,7 @@ private void logOnClientException(
* @see PipeConnector#transfer(TsFileInsertionEvent) for more details.
*/
private void transferQueuedEventsIfNecessary(final boolean forced) {
throwIfMissingFileRetryLimitExceeded();
throwIfReceiverProbeIsDelayed();

if ((retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty())
Expand Down Expand Up @@ -726,6 +747,7 @@ private void transferQueuedEventsIfNecessary(final boolean forced) {
}
}

throwIfMissingFileRetryLimitExceeded();
throwIfReceiverProbeIsDelayed();

// Stop retrying if the execution time exceeds the threshold for better realtime performance
Expand Down Expand Up @@ -868,17 +890,20 @@ private synchronized void addFailureEventToRetryQueue(
final EnrichedEvent enrichedEvent = (EnrichedEvent) event;
if (enrichedEvent.isReleased()) {
retryingEvent2Handlers.remove(event);
missingFileRetryTimes.remove(event);
return;
}
if (isDroppedPipe(enrichedEvent)) {
retryingEvent2Handlers.remove(event);
missingFileRetryTimes.remove(event);
enrichedEvent.clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
return;
}
}

if (isClosed.get()) {
retryingEvent2Handlers.remove(event);
missingFileRetryTimes.remove(event);
if (event instanceof EnrichedEvent) {
((EnrichedEvent) event).clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
}
Expand All @@ -887,6 +912,7 @@ private synchronized void addFailureEventToRetryQueue(

if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
lastRetryFailureMessage = null;
missingFileRetryLimitException = null;
}

if (e != null) {
Expand All @@ -908,6 +934,14 @@ private synchronized void addFailureEventToRetryQueue(
final PipeResourceFailureType previousResourceFailureType =
retryEvent2ResourceFailureType.get(event);

if (!alreadyInRetryQueue && isMissingFileFailure(event, e)) {
final int failureCount = missingFileRetryTimes.merge(event, 1, Integer::sum);
if (failureCount > MAX_FILE_NOT_FOUND_RETRY_TIMES) {
reportMissingFileEventAfterRetryLimit(event, e, failureCount - 1);
return;
}
}

if (resourceFailureType != null
&& event instanceof EnrichedEvent
&& (!alreadyInRetryQueue || previousResourceFailureType != resourceFailureType)) {
Expand Down Expand Up @@ -945,6 +979,64 @@ private synchronized void addFailureEventToRetryQueue(
}
}

private static boolean isMissingFileFailure(final Event event, final Exception exception) {
if (!(event instanceof PipeTsFileInsertionEvent) || exception == null) {
return false;
}
return ErrorHandlingCommonUtils.getRootCause(exception) instanceof FileNotFoundException;
}

private void reportMissingFileEventAfterRetryLimit(
final Event event, final Exception cause, final int retryTimes) {
final String missingFile = getMissingFile(cause);
final File tsFile = ((PipeTsFileInsertionEvent) event).getTsFile();
final String message =
String.format(
"Failed to transfer TsFile %s because file %s is missing after %s retries.",
tsFile, missingFile, retryTimes);
final PipeRuntimeSinkCriticalException criticalException =
new PipeRuntimeSinkCriticalException(message, cause);

retryingEvent2Handlers.remove(event);
retryEvent2ResourceFailureType.put(event, null);
missingFileRetryTimes.remove(event);
lastRetryFailureMessage = message;
missingFileRetryLimitException = criticalException;
retryTsFileQueue.offer((PipeTsFileInsertionEvent) event);
retryEventQueueEventCounter.increaseEventCount(event);
if (event instanceof EnrichedEvent) {
final EnrichedEvent enrichedEvent = (EnrichedEvent) event;
PipeLogger.log(
LOGGER::error,
criticalException,
"Failed to transfer TsFile %s because file %s is missing after %s retries.",
tsFile,
missingFile,
retryTimes);
PipeDataNodeAgent.runtime().report(enrichedEvent, criticalException);
} else {
LOGGER.error(message, cause);
}
}

private void throwIfMissingFileRetryLimitExceeded() {
final PipeRuntimeSinkCriticalException exception = missingFileRetryLimitException;
if (exception != null) {
throw exception;
}
}

private static String getMissingFile(final Exception exception) {
final Throwable rootCause = ErrorHandlingCommonUtils.getRootCause(exception);
return hasText(rootCause.getMessage()) ? rootCause.getMessage() : rootCause.toString();
}

public synchronized void clearFileNotFoundRetryTimes(final Iterable<? extends Event> events) {
for (final Event event : events) {
missingFileRetryTimes.remove(event);
}
}

/**
* Add failure {@link EnrichedEvent}s to retry queue.
*
Expand Down Expand Up @@ -1150,6 +1242,7 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {
((EnrichedEvent) event).clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
retryEventQueueEventCounter.decreaseEventCount(event);
retryEvent2ResourceFailureType.remove(event);
missingFileRetryTimes.remove(event);
retryingEvent2Handlers.remove(event);
return true;
}
Expand All @@ -1163,6 +1256,7 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {
((EnrichedEvent) event).clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
retryEventQueueEventCounter.decreaseEventCount(event);
retryEvent2ResourceFailureType.remove(event);
missingFileRetryTimes.remove(event);
retryingEvent2Handlers.remove(event);
return true;
}
Expand All @@ -1178,6 +1272,7 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {

if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
lastRetryFailureMessage = null;
missingFileRetryLimitException = null;
}
}

Expand Down Expand Up @@ -1230,12 +1325,15 @@ public synchronized void clearRetryEventsReferenceCount() {
retryTsFileQueue.isEmpty() ? retryEventQueue.poll() : retryTsFileQueue.poll();
retryEventQueueEventCounter.decreaseEventCount(event);
retryEvent2ResourceFailureType.remove(event);
missingFileRetryTimes.remove(event);
retryingEvent2Handlers.remove(event);
if (event instanceof EnrichedEvent) {
((EnrichedEvent) event).clearReferenceCount(IoTDBDataRegionAsyncSink.class.getName());
}
}
retryEvent2ResourceFailureType.clear();
missingFileRetryTimes.clear();
missingFileRetryLimitException = null;
lastRetryFailureMessage = null;
retryingEvent2Handlers.clear();
}
Expand Down
Loading
Loading