2424import org .apache .iotdb .commons .client .ThriftClient ;
2525import org .apache .iotdb .commons .client .async .AsyncPipeDataTransferServiceClient ;
2626import org .apache .iotdb .commons .exception .pipe .PipeRuntimeOutOfMemoryCriticalException ;
27+ import org .apache .iotdb .commons .exception .pipe .PipeRuntimeSinkCriticalException ;
2728import org .apache .iotdb .commons .exception .pipe .PipeRuntimeSinkNonReportTimeConfigurableException ;
2829import org .apache .iotdb .commons .exception .pipe .PipeRuntimeSinkResourceException ;
2930import org .apache .iotdb .commons .pipe .agent .task .progress .CommitterKey ;
@@ -119,6 +120,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
119120 "Failed to retry transferring events in the retry queue. Remaining events: %d (tablet events: %d, tsfile events: %d). Last failure: %s." ;
120121
121122 private static final boolean isSplitTSFileBatchModeEnabled = true ;
123+ private static final int MAX_FILE_NOT_FOUND_RETRY_TIMES = 5 ;
122124
123125 private final IoTDBDataRegionSyncSink syncConnector = new IoTDBDataRegionSyncSink ();
124126 private final BlockingQueue <Event > retryEventQueue = new LinkedBlockingQueue <>();
@@ -129,6 +131,10 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
129131 // semantics because the same payload may compare equal.
130132 private final Map <Event , PipeResourceFailureType > retryEvent2ResourceFailureType =
131133 new IdentityHashMap <>();
134+ // File disappearance is permanent for a given TsFile event. Keep its retry count separate from
135+ // transient receiver failures so one missing file cannot retry forever in the async queue.
136+ private final Map <Event , Integer > missingFileRetryTimes = new IdentityHashMap <>();
137+ private volatile PipeRuntimeSinkCriticalException missingFileRetryLimitException ;
132138 // Keep only the latest text to avoid retaining the complete exception chain for every event.
133139 private volatile String lastRetryFailureMessage ;
134140 // Events removed from the retry queue remain here while their next transfer is being started.
@@ -510,6 +516,21 @@ private boolean transferWithoutCheck(final TsFileInsertionEvent tsFileInsertionE
510516 throw new FileNotFoundException (pipeTsFileInsertionEvent .getTsFile ().getAbsolutePath ());
511517 }
512518
519+ final boolean transferMod =
520+ pipeTsFileInsertionEvent .isWithMod ()
521+ && clientManager .supportModsIfIsDataNodeReceiver ()
522+ && pipeTsFileInsertionEvent .getModFile () != null
523+ && pipeTsFileInsertionEvent .getModFile ().exists ();
524+ if (pipeTsFileInsertionEvent .isWithMod ()
525+ && clientManager .supportModsIfIsDataNodeReceiver ()
526+ && pipeTsFileInsertionEvent .getModFile () != null
527+ && !pipeTsFileInsertionEvent .getModFile ().exists ()) {
528+ LOGGER .warn (
529+ "TsFile {}: modification file {} is missing, transfer the TsFile without modifications." ,
530+ pipeTsFileInsertionEvent .getTsFile (),
531+ pipeTsFileInsertionEvent .getModFile ());
532+ }
533+
513534 pipeTransferTsFileHandler =
514535 new PipeTransferTsFileHandler (
515536 this ,
@@ -523,8 +544,7 @@ private boolean transferWithoutCheck(final TsFileInsertionEvent tsFileInsertionE
523544 new AtomicBoolean (false ),
524545 pipeTsFileInsertionEvent .getTsFile (),
525546 pipeTsFileInsertionEvent .getModFile (),
526- pipeTsFileInsertionEvent .isWithMod ()
527- && clientManager .supportModsIfIsDataNodeReceiver (),
547+ transferMod ,
528548 pipeTsFileInsertionEvent .getDatabaseName ());
529549 trackRetryHandler (
530550 pipeTransferTsFileHandler , Collections .singletonList (pipeTsFileInsertionEvent ));
@@ -661,6 +681,7 @@ private void logOnClientException(
661681 * @see PipeConnector#transfer(TsFileInsertionEvent) for more details.
662682 */
663683 private void transferQueuedEventsIfNecessary (final boolean forced ) {
684+ throwIfMissingFileRetryLimitExceeded ();
664685 throwIfReceiverProbeIsDelayed ();
665686
666687 if ((retryEventQueue .isEmpty () && retryTsFileQueue .isEmpty ())
@@ -726,6 +747,7 @@ private void transferQueuedEventsIfNecessary(final boolean forced) {
726747 }
727748 }
728749
750+ throwIfMissingFileRetryLimitExceeded ();
729751 throwIfReceiverProbeIsDelayed ();
730752
731753 // Stop retrying if the execution time exceeds the threshold for better realtime performance
@@ -868,17 +890,20 @@ private synchronized void addFailureEventToRetryQueue(
868890 final EnrichedEvent enrichedEvent = (EnrichedEvent ) event ;
869891 if (enrichedEvent .isReleased ()) {
870892 retryingEvent2Handlers .remove (event );
893+ missingFileRetryTimes .remove (event );
871894 return ;
872895 }
873896 if (isDroppedPipe (enrichedEvent )) {
874897 retryingEvent2Handlers .remove (event );
898+ missingFileRetryTimes .remove (event );
875899 enrichedEvent .clearReferenceCount (IoTDBDataRegionAsyncSink .class .getName ());
876900 return ;
877901 }
878902 }
879903
880904 if (isClosed .get ()) {
881905 retryingEvent2Handlers .remove (event );
906+ missingFileRetryTimes .remove (event );
882907 if (event instanceof EnrichedEvent ) {
883908 ((EnrichedEvent ) event ).clearReferenceCount (IoTDBDataRegionAsyncSink .class .getName ());
884909 }
@@ -887,6 +912,7 @@ private synchronized void addFailureEventToRetryQueue(
887912
888913 if (retryEventQueue .isEmpty () && retryTsFileQueue .isEmpty ()) {
889914 lastRetryFailureMessage = null ;
915+ missingFileRetryLimitException = null ;
890916 }
891917
892918 if (e != null ) {
@@ -908,6 +934,14 @@ private synchronized void addFailureEventToRetryQueue(
908934 final PipeResourceFailureType previousResourceFailureType =
909935 retryEvent2ResourceFailureType .get (event );
910936
937+ if (!alreadyInRetryQueue && isMissingFileFailure (event , e )) {
938+ final int failureCount = missingFileRetryTimes .merge (event , 1 , Integer ::sum );
939+ if (failureCount > MAX_FILE_NOT_FOUND_RETRY_TIMES ) {
940+ reportMissingFileEventAfterRetryLimit (event , e , failureCount - 1 );
941+ return ;
942+ }
943+ }
944+
911945 if (resourceFailureType != null
912946 && event instanceof EnrichedEvent
913947 && (!alreadyInRetryQueue || previousResourceFailureType != resourceFailureType )) {
@@ -945,6 +979,64 @@ private synchronized void addFailureEventToRetryQueue(
945979 }
946980 }
947981
982+ private static boolean isMissingFileFailure (final Event event , final Exception exception ) {
983+ if (!(event instanceof PipeTsFileInsertionEvent ) || exception == null ) {
984+ return false ;
985+ }
986+ return ErrorHandlingCommonUtils .getRootCause (exception ) instanceof FileNotFoundException ;
987+ }
988+
989+ private void reportMissingFileEventAfterRetryLimit (
990+ final Event event , final Exception cause , final int retryTimes ) {
991+ final String missingFile = getMissingFile (cause );
992+ final File tsFile = ((PipeTsFileInsertionEvent ) event ).getTsFile ();
993+ final String message =
994+ String .format (
995+ "Failed to transfer TsFile %s because file %s is missing after %s retries." ,
996+ tsFile , missingFile , retryTimes );
997+ final PipeRuntimeSinkCriticalException criticalException =
998+ new PipeRuntimeSinkCriticalException (message , cause );
999+
1000+ retryingEvent2Handlers .remove (event );
1001+ retryEvent2ResourceFailureType .put (event , null );
1002+ missingFileRetryTimes .remove (event );
1003+ lastRetryFailureMessage = message ;
1004+ missingFileRetryLimitException = criticalException ;
1005+ retryTsFileQueue .offer ((PipeTsFileInsertionEvent ) event );
1006+ retryEventQueueEventCounter .increaseEventCount (event );
1007+ if (event instanceof EnrichedEvent ) {
1008+ final EnrichedEvent enrichedEvent = (EnrichedEvent ) event ;
1009+ PipeLogger .log (
1010+ LOGGER ::error ,
1011+ criticalException ,
1012+ "Failed to transfer TsFile %s because file %s is missing after %s retries." ,
1013+ tsFile ,
1014+ missingFile ,
1015+ retryTimes );
1016+ PipeDataNodeAgent .runtime ().report (enrichedEvent , criticalException );
1017+ } else {
1018+ LOGGER .error (message , cause );
1019+ }
1020+ }
1021+
1022+ private void throwIfMissingFileRetryLimitExceeded () {
1023+ final PipeRuntimeSinkCriticalException exception = missingFileRetryLimitException ;
1024+ if (exception != null ) {
1025+ throw exception ;
1026+ }
1027+ }
1028+
1029+ private static String getMissingFile (final Exception exception ) {
1030+ final Throwable rootCause = ErrorHandlingCommonUtils .getRootCause (exception );
1031+ return hasText (rootCause .getMessage ()) ? rootCause .getMessage () : rootCause .toString ();
1032+ }
1033+
1034+ public synchronized void clearFileNotFoundRetryTimes (final Iterable <? extends Event > events ) {
1035+ for (final Event event : events ) {
1036+ missingFileRetryTimes .remove (event );
1037+ }
1038+ }
1039+
9481040 /**
9491041 * Add failure {@link EnrichedEvent}s to retry queue.
9501042 *
@@ -1150,6 +1242,7 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {
11501242 ((EnrichedEvent ) event ).clearReferenceCount (IoTDBDataRegionAsyncSink .class .getName ());
11511243 retryEventQueueEventCounter .decreaseEventCount (event );
11521244 retryEvent2ResourceFailureType .remove (event );
1245+ missingFileRetryTimes .remove (event );
11531246 retryingEvent2Handlers .remove (event );
11541247 return true ;
11551248 }
@@ -1163,6 +1256,7 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {
11631256 ((EnrichedEvent ) event ).clearReferenceCount (IoTDBDataRegionAsyncSink .class .getName ());
11641257 retryEventQueueEventCounter .decreaseEventCount (event );
11651258 retryEvent2ResourceFailureType .remove (event );
1259+ missingFileRetryTimes .remove (event );
11661260 retryingEvent2Handlers .remove (event );
11671261 return true ;
11681262 }
@@ -1178,6 +1272,7 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {
11781272
11791273 if (retryEventQueue .isEmpty () && retryTsFileQueue .isEmpty ()) {
11801274 lastRetryFailureMessage = null ;
1275+ missingFileRetryLimitException = null ;
11811276 }
11821277 }
11831278
@@ -1230,12 +1325,15 @@ public synchronized void clearRetryEventsReferenceCount() {
12301325 retryTsFileQueue .isEmpty () ? retryEventQueue .poll () : retryTsFileQueue .poll ();
12311326 retryEventQueueEventCounter .decreaseEventCount (event );
12321327 retryEvent2ResourceFailureType .remove (event );
1328+ missingFileRetryTimes .remove (event );
12331329 retryingEvent2Handlers .remove (event );
12341330 if (event instanceof EnrichedEvent ) {
12351331 ((EnrichedEvent ) event ).clearReferenceCount (IoTDBDataRegionAsyncSink .class .getName ());
12361332 }
12371333 }
12381334 retryEvent2ResourceFailureType .clear ();
1335+ missingFileRetryTimes .clear ();
1336+ missingFileRetryLimitException = null ;
12391337 lastRetryFailureMessage = null ;
12401338 retryingEvent2Handlers .clear ();
12411339 }
0 commit comments