This is an automated email from the ASF dual-hosted git repository.

jt2594838 pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/dev/1.3 by this push:
     new a9f69c34cb1 [Pipe] Back off async sink for all retry statuses (#18236) 
(#18252)
a9f69c34cb1 is described below

commit a9f69c34cb18f570253091c99b73b7cdeae84969
Author: Caideyipi <[email protected]>
AuthorDate: Tue Jul 21 09:55:34 2026 +0800

    [Pipe] Back off async sink for all retry statuses (#18236) (#18252)
---
 .../thrift/async/IoTDBDataRegionAsyncSink.java     | 41 ++++++++++------------
 .../handler/PipeTransferTrackableHandler.java      |  4 +--
 .../handler/PipeTransferTrackableHandlerTest.java  | 20 ++++++-----
 3 files changed, 31 insertions(+), 34 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
index 13589527b37..3d5d6af26a0 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
@@ -128,8 +128,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
       new ConcurrentHashMap<>();
 
   private final Set<CommitterKey> droppedPipeTaskKeys = 
ConcurrentHashMap.newKeySet();
-  private final Map<String, ReceiverTemporaryUnavailableBackoff> 
receiverBackoffMap =
-      new ConcurrentHashMap<>();
+  private final Map<String, ReceiverRetryBackoff> receiverBackoffMap = new 
ConcurrentHashMap<>();
 
   private boolean enableSendTsFileLimit;
   private volatile boolean isConnectionException;
@@ -740,13 +739,13 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
     return enableSendTsFileLimit;
   }
 
-  public void waitIfReceiverTemporarilyUnavailable(final TEndPoint endPoint) {
+  public void waitIfReceiverRetryIsBackedOff(final TEndPoint endPoint) {
     final String endPointKey = format(endPoint);
     if (Objects.isNull(endPointKey)) {
       return;
     }
 
-    final ReceiverTemporaryUnavailableBackoff backoff = 
receiverBackoffMap.get(endPointKey);
+    final ReceiverRetryBackoff backoff = receiverBackoffMap.get(endPointKey);
     if (Objects.isNull(backoff)) {
       return;
     }
@@ -782,8 +781,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
   }
 
   private void throwIfReceiverProbeIsDelayed() {
-    for (final Map.Entry<String, ReceiverTemporaryUnavailableBackoff> entry :
-        receiverBackoffMap.entrySet()) {
+    for (final Map.Entry<String, ReceiverRetryBackoff> entry : 
receiverBackoffMap.entrySet()) {
       final long probeDelayInMs = 
entry.getValue().getRemainingProbeDelayInMs();
       if (probeDelayInMs <= 0) {
         continue;
@@ -795,10 +793,10 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
 
   private static PipeRuntimeSinkNonReportTimeConfigurableException
       createReceiverProbeDelayException(
-          final String endPointKey, final ReceiverTemporaryUnavailableBackoff 
backoff) {
+          final String endPointKey, final ReceiverRetryBackoff backoff) {
     return new PipeRuntimeSinkNonReportTimeConfigurableException(
         String.format(
-            "Receiver %s remained temporarily unavailable for more than %d ms, 
pause regular retries and probe every %d ms.",
+            "Receiver %s has required retries for more than %d ms, pause 
regular retries and probe every %d ms.",
             endPointKey, backoff.getRetryMaxDurationInMs(), 
backoff.getRetryProbeIntervalInMs()),
         Long.MAX_VALUE);
   }
@@ -809,14 +807,14 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
       return;
     }
 
-    if (isReceiverTemporarilyUnavailable(status)) {
+    if (isReceiverRetryNeeded(status)) {
       final long backoffTimeInMs =
           receiverBackoffMap
-              .computeIfAbsent(endPointKey, key -> new 
ReceiverTemporaryUnavailableBackoff())
-              .markTemporarilyUnavailable();
+              .computeIfAbsent(endPointKey, key -> new ReceiverRetryBackoff())
+              .markRetryNeeded();
       if (LOGGER.isDebugEnabled()) {
         LOGGER.debug(
-            "Receiver {} is temporarily unavailable, throttle requests for {} 
ms. Status: {}",
+            "Receiver {} requires a retry, throttle requests for {} ms. 
Status: {}",
             endPointKey,
             backoffTimeInMs,
             status);
@@ -834,20 +832,17 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
     }
   }
 
-  private static boolean isReceiverTemporarilyUnavailable(final TSStatus 
status) {
+  private static boolean isReceiverRetryNeeded(final TSStatus status) {
     if (Objects.isNull(status)) {
       return false;
     }
 
-    final int statusCode = status.getCode();
-    if (statusCode == 
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()
-        || statusCode == TSStatusCode.WRITE_PROCESS_REJECT.getStatusCode()) {
+    if (!isSuccess(status)) {
       return true;
     }
 
     return status.isSetSubStatus()
-        && status.getSubStatus().stream()
-            
.anyMatch(IoTDBDataRegionAsyncSink::isReceiverTemporarilyUnavailable);
+        && 
status.getSubStatus().stream().anyMatch(IoTDBDataRegionAsyncSink::isReceiverRetryNeeded);
   }
 
   private static boolean isSuccess(final TSStatus status) {
@@ -1036,7 +1031,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
     }
   }
 
-  private static class ReceiverTemporaryUnavailableBackoff {
+  private static class ReceiverRetryBackoff {
 
     private final long maxBackoffTimeInMs =
         Math.max(0, 
PipeConfig.getInstance().getPipeSinkSubtaskSleepIntervalMaxMs());
@@ -1050,17 +1045,17 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
{
         Math.max(1, 
PipeConfig.getInstance().getPipeAsyncSinkRetryProbeIntervalMs());
 
     private boolean active = false;
-    private long firstUnavailableTimeInMs = 0;
+    private long firstRetryTimeInMs = 0;
     private long currentBackoffTimeInMs = initialBackoffTimeInMs;
     private long failureBackoffUntilInMs = 0;
     private long nextReservedRetryTimeInMs = 0;
     private long nextProbeTimeInMs = 0;
 
-    private synchronized long markTemporarilyUnavailable() {
+    private synchronized long markRetryNeeded() {
       final long currentTimeInMs = System.currentTimeMillis();
       if (!active) {
         active = true;
-        firstUnavailableTimeInMs = currentTimeInMs;
+        firstRetryTimeInMs = currentTimeInMs;
         currentBackoffTimeInMs = initialBackoffTimeInMs;
         failureBackoffUntilInMs = 0;
         nextReservedRetryTimeInMs = 0;
@@ -1082,7 +1077,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
     private synchronized boolean isRetryMaxDurationExceeded() {
       return active
           && retryMaxDurationInMs >= 0
-          && System.currentTimeMillis() - firstUnavailableTimeInMs >= 
retryMaxDurationInMs;
+          && System.currentTimeMillis() - firstRetryTimeInMs >= 
retryMaxDurationInMs;
     }
 
     private synchronized long reserveNextRetryTimeInMs() {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandler.java
index edae06ac65f..b2c91495496 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandler.java
@@ -108,7 +108,7 @@ public abstract class PipeTransferTrackableHandler
       return false;
     }
     try {
-      sink.waitIfReceiverTemporarilyUnavailable(client.getEndPoint());
+      sink.waitIfReceiverRetryIsBackedOff(client.getEndPoint());
     } catch (final PipeRuntimeSinkNonReportTimeConfigurableException e) {
       returnClientToPool(client);
       onError(e);
@@ -293,7 +293,7 @@ public abstract class PipeTransferTrackableHandler
 
     try {
       client.setShouldReturnSelf(shouldReturnSelf);
-      sink.waitIfReceiverTemporarilyUnavailable(client.getEndPoint());
+      sink.waitIfReceiverRetryIsBackedOff(client.getEndPoint());
       if (returnFalseIfSinkIsClosed(client)) {
         return;
       }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandlerTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandlerTest.java
index 61218d112c8..a0f71536fff 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandlerTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTrackableHandlerTest.java
@@ -194,7 +194,7 @@ public class PipeTransferTrackableHandlerTest {
     handler.transfer(client, createReq(1));
 
     final InOrder inOrder = Mockito.inOrder(sink, client);
-    inOrder.verify(sink).waitIfReceiverTemporarilyUnavailable(endPoint);
+    inOrder.verify(sink).waitIfReceiverRetryIsBackedOff(endPoint);
     inOrder.verify(client).pipeTransfer(Mockito.any(TPipeTransferReq.class), 
Mockito.any());
     Mockito.verify(sink).recordReceiverStatus(endPoint, status);
   }
@@ -208,7 +208,7 @@ public class PipeTransferTrackableHandlerTest {
     final PipeRuntimeSinkNonReportTimeConfigurableException exception =
         new PipeRuntimeSinkNonReportTimeConfigurableException("probe delayed", 
Long.MAX_VALUE);
     Mockito.when(client.getEndPoint()).thenReturn(endPoint);
-    
Mockito.doThrow(exception).when(sink).waitIfReceiverTemporarilyUnavailable(endPoint);
+    
Mockito.doThrow(exception).when(sink).waitIfReceiverRetryIsBackedOff(endPoint);
 
     final TestPipeTransferTrackableHandler handler = new 
TestPipeTransferTrackableHandler(sink);
 
@@ -221,18 +221,19 @@ public class PipeTransferTrackableHandlerTest {
   }
 
   @Test
-  public void testReceiverRetriesAreSerialized() {
+  public void testReceiverRetriesAreSerializedForAnyFailureStatus() {
     commonConfig.setPipeSinkSubtaskSleepIntervalInitMs(40);
     commonConfig.setPipeSinkSubtaskSleepIntervalMaxMs(40);
     commonConfig.setPipeAsyncSinkRetryMaxDurationMs(5000);
 
     final IoTDBDataRegionAsyncSink sink = new IoTDBDataRegionAsyncSink();
     final TEndPoint endPoint = new TEndPoint("127.0.0.1", 6667);
-    sink.recordReceiverStatus(endPoint, temporarilyUnavailableStatus());
+    sink.recordReceiverStatus(
+        endPoint, new 
TSStatus().setCode(TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode()));
 
     final long startTimeInMs = System.currentTimeMillis();
-    sink.waitIfReceiverTemporarilyUnavailable(endPoint);
-    sink.waitIfReceiverTemporarilyUnavailable(endPoint);
+    sink.waitIfReceiverRetryIsBackedOff(endPoint);
+    sink.waitIfReceiverRetryIsBackedOff(endPoint);
 
     Assert.assertTrue(System.currentTimeMillis() - startTimeInMs >= 60);
   }
@@ -246,13 +247,14 @@ public class PipeTransferTrackableHandlerTest {
     final TEndPoint endPoint = new TEndPoint("127.0.0.1", 6667);
     sink.recordReceiverStatus(endPoint, temporarilyUnavailableStatus());
 
-    sink.waitIfReceiverTemporarilyUnavailable(endPoint);
+    sink.waitIfReceiverRetryIsBackedOff(endPoint);
     Assert.assertThrows(
         PipeRuntimeSinkNonReportTimeConfigurableException.class,
-        () -> sink.waitIfReceiverTemporarilyUnavailable(endPoint));
+        () -> sink.waitIfReceiverRetryIsBackedOff(endPoint));
+
     sink.recordReceiverStatus(
         endPoint, new 
TSStatus().setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
-    sink.waitIfReceiverTemporarilyUnavailable(endPoint);
+    sink.waitIfReceiverRetryIsBackedOff(endPoint);
   }
 
   private static TSStatus temporarilyUnavailableStatus() {

Reply via email to