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

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


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

commit 94b07e08fdaee976c961e4a7bf2f8f6fa1d57d7f
Author: Caideyipi <[email protected]>
AuthorDate: Mon Jul 20 15:12:29 2026 +0800

    [Pipe] Back off async sink for all retry statuses (#18236)
---
 .../apache/iotdb/db/i18n/DataNodePipeMessages.java |  8 ++---
 .../apache/iotdb/db/i18n/DataNodePipeMessages.java |  8 ++---
 .../thrift/async/IoTDBDataRegionAsyncSink.java     | 41 ++++++++++------------
 .../handler/PipeTransferTrackableHandler.java      |  4 +--
 .../handler/PipeTransferTrackableHandlerTest.java  | 19 +++++-----
 5 files changed, 38 insertions(+), 42 deletions(-)

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 bfd57672ee3..40b499410aa 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
@@ -2541,10 +2541,10 @@ public final class DataNodePipeMessages {
       "Incomplete column values in current tablet format deserialization.";
   public static final String 
EXCEPTION_INCOMPLETE_TIMESTAMPS_IN_CURRENT_TABLET_FORMAT_DESERIALIZATION_FE212461
 =
       "Incomplete timestamps in current tablet format deserialization.";
-  public static final String 
MESSAGE_RECEIVER_ARG_IS_TEMPORARILY_UNAVAILABLE_THROTTLE_REQUESTS_FOR_ARG_MS_STATUS_ARG_F37192D9
 =
-      "Receiver {} is temporarily unavailable, throttle requests for {} ms. 
Status: {}";
-  public static final String 
EXCEPTION_RECEIVER_ARG_REMAINED_TEMPORARILY_UNAVAILABLE_FOR_MORE_THAN_ARG_MS_PAUSE_REGULAR_RETRIES_AND_PROBE_EVERY_ARG_MS_C515DD97
 =
-      "Receiver %s remained temporarily unavailable for more than %d ms, pause 
regular retries and probe every %d ms.";
+  public static final String 
MESSAGE_RECEIVER_ARG_REQUIRES_A_RETRY_THROTTLE_REQUESTS_FOR_ARG_MS_STATUS_ARG_0B3B14F6
 =
+      "Receiver {} requires a retry, throttle requests for {} ms. Status: {}";
+  public static final String 
EXCEPTION_RECEIVER_ARG_HAS_REQUIRED_RETRIES_FOR_MORE_THAN_ARG_MS_PAUSE_REGULAR_RETRIES_AND_PROBE_EVERY_ARG_MS_550475C2
 =
+      "Receiver %s has required retries for more than %d ms, pause regular 
retries and probe every %d ms.";
   public static final String 
MESSAGE_SUCCESSFULLY_TRANSFERRED_BATCHED_SCHEMA_EVENTS_BATCH_SIZE_ARG_CF2E881C =
       "Successfully transferred batched schema events, batch size {}.";
   public static final String 
EXCEPTION_AUTO_CREATE_TREE_DATABASE_FAILED_ARG_STATUS_CODE_ARG_C6175C27 =
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 2b97336eaba..85a8ae289cf 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
@@ -2370,10 +2370,10 @@ public final class DataNodePipeMessages {
       "当前 tablet 格式反序列化中列值不完整。";
   public static final String 
EXCEPTION_INCOMPLETE_TIMESTAMPS_IN_CURRENT_TABLET_FORMAT_DESERIALIZATION_FE212461
 =
       "当前 tablet 格式反序列化中时间戳不完整。";
-  public static final String 
MESSAGE_RECEIVER_ARG_IS_TEMPORARILY_UNAVAILABLE_THROTTLE_REQUESTS_FOR_ARG_MS_STATUS_ARG_F37192D9
 =
-      "Receiver {} 暂时不可用,对请求限流 {} ms。状态:{}";
-  public static final String 
EXCEPTION_RECEIVER_ARG_REMAINED_TEMPORARILY_UNAVAILABLE_FOR_MORE_THAN_ARG_MS_PAUSE_REGULAR_RETRIES_AND_PROBE_EVERY_ARG_MS_C515DD97
 =
-      "Receiver %s 持续暂时不可用超过 %d ms,暂停常规重试,改为每 %d ms 探测一次。";
+  public static final String 
MESSAGE_RECEIVER_ARG_REQUIRES_A_RETRY_THROTTLE_REQUESTS_FOR_ARG_MS_STATUS_ARG_0B3B14F6
 =
+      "Receiver {} 要求重试,对请求限流 {} ms。状态:{}";
+  public static final String 
EXCEPTION_RECEIVER_ARG_HAS_REQUIRED_RETRIES_FOR_MORE_THAN_ARG_MS_PAUSE_REGULAR_RETRIES_AND_PROBE_EVERY_ARG_MS_550475C2
 =
+      "Receiver %s 要求重试已超过 %d ms,暂停常规重试,改为每 %d ms 探测一次。";
   public static final String 
MESSAGE_SUCCESSFULLY_TRANSFERRED_BATCHED_SCHEMA_EVENTS_BATCH_SIZE_ARG_CF2E881C =
       "成功传输批量的 schema 事件,batch 大小 {}。";
   public static final String 
EXCEPTION_AUTO_CREATE_TREE_DATABASE_FAILED_ARG_STATUS_CODE_ARG_C6175C27 =
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 6153d19e444..480ade16546 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
@@ -147,8 +147,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
       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;
@@ -843,13 +842,13 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
     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;
     }
@@ -886,8 +885,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
   }
 
   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;
@@ -900,11 +898,11 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
 
   private static PipeRuntimeSinkNonReportTimeConfigurableException
       createReceiverProbeDelayException(
-          final String endPointKey, final ReceiverTemporaryUnavailableBackoff 
backoff) {
+          final String endPointKey, final ReceiverRetryBackoff backoff) {
     return new PipeRuntimeSinkNonReportTimeConfigurableException(
         String.format(
             DataNodePipeMessages
-                
.EXCEPTION_RECEIVER_ARG_REMAINED_TEMPORARILY_UNAVAILABLE_FOR_MORE_THAN_ARG_MS_PAUSE_REGULAR_RETRIES_AND_PROBE_EVERY_ARG_MS_C515DD97,
+                
.EXCEPTION_RECEIVER_ARG_HAS_REQUIRED_RETRIES_FOR_MORE_THAN_ARG_MS_PAUSE_REGULAR_RETRIES_AND_PROBE_EVERY_ARG_MS_550475C2,
             endPointKey,
             backoff.getRetryMaxDurationInMs(),
             backoff.getRetryProbeIntervalInMs()),
@@ -917,15 +915,15 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
       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(
             DataNodePipeMessages
-                
.MESSAGE_RECEIVER_ARG_IS_TEMPORARILY_UNAVAILABLE_THROTTLE_REQUESTS_FOR_ARG_MS_STATUS_ARG_F37192D9,
+                
.MESSAGE_RECEIVER_ARG_REQUIRES_A_RETRY_THROTTLE_REQUESTS_FOR_ARG_MS_STATUS_ARG_0B3B14F6,
             endPointKey,
             backoffTimeInMs,
             status);
@@ -946,20 +944,17 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
     }
   }
 
-  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) {
@@ -1148,7 +1143,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
     }
   }
 
-  private static class ReceiverTemporaryUnavailableBackoff {
+  private static class ReceiverRetryBackoff {
 
     private final long maxBackoffTimeInMs =
         Math.max(0, 
PipeConfig.getInstance().getPipeSinkSubtaskSleepIntervalMaxMs());
@@ -1162,17 +1157,17 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
         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;
@@ -1194,7 +1189,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink 
implements PipeSinkWithS
     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 a3792bcb3da..e52c2479c04 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
@@ -109,7 +109,7 @@ public abstract class PipeTransferTrackableHandler
       return false;
     }
     try {
-      sink.waitIfReceiverTemporarilyUnavailable(client.getEndPoint());
+      sink.waitIfReceiverRetryIsBackedOff(client.getEndPoint());
     } catch (final PipeRuntimeSinkNonReportTimeConfigurableException e) {
       returnClientToPool(client);
       onError(e);
@@ -299,7 +299,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 49d1f617afa..28f4f910b0f 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,15 +247,15 @@ 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));
     Assert.assertTrue(sink.peekSchedulingDelayMs() > 0);
 
     sink.recordReceiverStatus(
         endPoint, new 
TSStatus().setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
-    sink.waitIfReceiverTemporarilyUnavailable(endPoint);
+    sink.waitIfReceiverRetryIsBackedOff(endPoint);
   }
 
   private static TSStatus temporarilyUnavailableStatus() {

Reply via email to