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 18ec5f0a393 [Pipe] Preserve downstream sink error messages (#18553)
18ec5f0a393 is described below
commit 18ec5f0a393ed7cee0e3d30df5bf96e034d565d8
Author: Caideyipi <[email protected]>
AuthorDate: Wed Sep 2 10:13:30 2026 +0800
[Pipe] Preserve downstream sink error messages (#18553)
---
.../apache/iotdb/db/i18n/DataNodePipeMessages.java | 7 +
.../apache/iotdb/db/i18n/DataNodePipeMessages.java | 7 +
.../thrift/async/IoTDBDataRegionAsyncSink.java | 74 ++++++-
.../receiver/PipeStatementTsStatusVisitorTest.java | 18 ++
.../thrift/async/IoTDBDataRegionAsyncSinkTest.java | 66 ++++++
.../PipeTransferTsFileHandlerCleanupTest.java | 56 +++++
.../pipe/receiver/PipeReceiverStatusHandler.java | 243 ++++++++++++++++++++-
.../PipeReceiverStatusHandlerTest.java | 226 +++++++++++++++++++
8 files changed, 681 insertions(+), 16 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 41f5c2831bb..520a300c8e1 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
@@ -2590,6 +2590,13 @@ public final class DataNodePipeMessages {
public static final String
MESSAGE_TRANSFER_FILE_ARG_ERROR_RESULT_STATUS_ARG_E565D9FD =
"Transfer file %s error, result status %s.";
+ public static final String
+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_5B4B2E7C
=
+ "Failed to retry transferring events in the retry queue. Remaining
events: %d (tablet events: %d, tsfile events: %d).";
+ public static final String
+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_LAST_FAILURE_ARG_EB8F9DCD
=
+ "Failed to retry transferring events in the retry queue. Remaining
events: %d (tablet events: %d, tsfile events: %d). Last failure: %s.";
+
public static final String
EXCEPTION_LEGACY_PIPE_RECEIVER_REQUIRES_A_LOGGED_IN_SESSION_D96219BF =
"Legacy pipe receiver requires a logged-in session.";
public static final String
EXCEPTION_FAILED_TO_SET_UP_CONSENSUS_SUBSCRIPTION_FOR_TOPIC_ARG_IN_CONSUMER_GROUP_ARG_ARG_A7FA88F3
=
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 5a4d9c4c568..2a43385342f 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
@@ -2416,6 +2416,13 @@ public final class DataNodePipeMessages {
public static final String
MESSAGE_TRANSFER_FILE_ARG_ERROR_RESULT_STATUS_ARG_E565D9FD =
"传输文件 %s 出错,结果状态为 %s。";
+ public static final String
+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_5B4B2E7C
=
+ "重试 retry queue 中的事件失败。剩余事件:%d(tablet 事件:%d,tsfile 事件:%d)。";
+ public static final String
+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_LAST_FAILURE_ARG_EB8F9DCD
=
+ "重试 retry queue 中的事件失败。剩余事件:%d(tablet 事件:%d,tsfile
事件:%d)。最近一次失败:%s。";
+
public static final String
EXCEPTION_LEGACY_PIPE_RECEIVER_REQUIRES_A_LOGGED_IN_SESSION_D96219BF =
"Legacy pipe receiver 需要已登录的 session。";
public static final String
EXCEPTION_FAILED_TO_SET_UP_CONSENSUS_SUBSCRIPTION_FOR_TOPIC_ARG_IN_CONSUMER_GROUP_ARG_ARG_A7FA88F3
=
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 56e718c37a7..2102f743059 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
@@ -34,6 +34,7 @@ import
org.apache.iotdb.commons.pipe.resource.PipeStopStrategy;
import org.apache.iotdb.commons.pipe.resource.log.PipeLogger;
import org.apache.iotdb.commons.pipe.sink.protocol.IoTDBSink;
import org.apache.iotdb.commons.pipe.sink.protocol.PipeSinkWithSchedulingDelay;
+import org.apache.iotdb.commons.utils.ErrorHandlingCommonUtils;
import org.apache.iotdb.db.i18n.DataNodePipeMessages;
import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
import org.apache.iotdb.db.pipe.event.common.deletion.PipeDeleteDataNodeEvent;
@@ -138,6 +139,8 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
// Guarded by this. Events need identity semantics because the same payload
may compare equal.
private final Map<Event, PipeResourceFailureType>
retryEvent2ResourceFailureType =
new IdentityHashMap<>();
+ // Keep only the latest text to avoid retaining the complete exception chain
for every event.
+ private volatile String lastRetryFailureMessage;
private IoTDBDataNodeAsyncClientManager clientManager;
private IoTDBDataNodeAsyncClientManager transferTsFileClientManager;
@@ -732,13 +735,11 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
if (remainingEvents <= retryEventQueue.size() +
retryTsFileQueue.size()) {
final String message =
- "Failed to retry transferring events in the retry queue.
Remaining events: "
- + (retryEventQueue.size() + retryTsFileQueue.size())
- + " (tablet events: "
- + retryEventQueueEventCounter.getTabletInsertionEventCount()
- + ", tsfile events: "
- + retryEventQueueEventCounter.getTsFileInsertionEventCount()
- + ").";
+ formatRetryQueueFailureMessage(
+ retryEventQueue.size() + retryTsFileQueue.size(),
+ retryEventQueueEventCounter.getTabletInsertionEventCount(),
+ retryEventQueueEventCounter.getTsFileInsertionEventCount(),
+ lastRetryFailureMessage);
final PipeResourceFailureType retryQueueResourceFailureType =
getRetryQueueResourceFailureType();
if (retryQueueResourceFailureType != null) {
@@ -751,6 +752,12 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
}
}
}
+
+ synchronized (this) {
+ if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
+ lastRetryFailureMessage = null;
+ }
+ }
}
private void retryTransfer(final TabletInsertionEvent tabletInsertionEvent) {
@@ -834,6 +841,14 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
return;
}
+ if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
+ lastRetryFailureMessage = null;
+ }
+
+ if (e != null) {
+ lastRetryFailureMessage = getRetryFailureMessage(e);
+ }
+
if (resourceFailureType != null && event instanceof EnrichedEvent) {
final EnrichedEvent enrichedEvent = (EnrichedEvent) event;
final Pair<String, Long> pipeKey =
@@ -881,6 +896,46 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
events.forEach(event -> addFailureEventToRetryQueue(event, e,
failureRecordedPipes));
}
+ static String formatRetryQueueFailureMessage(
+ final int remainingEvents,
+ final int tabletEventCount,
+ final int tsFileEventCount,
+ final String lastFailureMessage) {
+ if (!hasText(lastFailureMessage)) {
+ return String.format(
+ DataNodePipeMessages
+
.EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_5B4B2E7C,
+ remainingEvents,
+ tabletEventCount,
+ tsFileEventCount);
+ }
+ return String.format(
+ DataNodePipeMessages
+
.EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_LAST_FAILURE_ARG_EB8F9DCD,
+ remainingEvents,
+ tabletEventCount,
+ tsFileEventCount,
+ lastFailureMessage);
+ }
+
+ private static String getRetryFailureMessage(final Exception exception) {
+ final Throwable rootCause =
ErrorHandlingCommonUtils.getRootCause(exception);
+ if (hasText(rootCause.getMessage())) {
+ return rootCause.getMessage();
+ }
+ // Throwable#getMessage() is null for exceptions such as a bare NPE. Keep
the type in the
+ // reported sink error instead of falling back to a generic transfer
wrapper.
+ return rootCause.toString();
+ }
+
+ private static boolean hasText(final String message) {
+ return message != null && !message.trim().isEmpty();
+ }
+
+ synchronized String getLastRetryFailureMessage() {
+ return lastRetryFailureMessage;
+ }
+
private synchronized PipeResourceFailureType
getRetryQueueResourceFailureType() {
for (final PipeResourceFailureType failureType :
PipeResourceFailureType.values()) {
if (retryEvent2ResourceFailureType.containsValue(failureType)) {
@@ -1056,6 +1111,10 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
}
return false;
});
+
+ if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
+ lastRetryFailureMessage = null;
+ }
}
@Override
@@ -1110,6 +1169,7 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink
implements PipeSinkWithS
}
}
retryEvent2ResourceFailureType.clear();
+ lastRetryFailureMessage = null;
}
//////////////////////// APIs provided for metric framework
////////////////////////
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/PipeStatementTsStatusVisitorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/PipeStatementTsStatusVisitorTest.java
index f2716d5c1a4..7bc1a60135e 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/PipeStatementTsStatusVisitorTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/PipeStatementTsStatusVisitorTest.java
@@ -66,6 +66,24 @@ public class PipeStatementTsStatusVisitorTest {
.getCode());
}
+ @Test
+ public void testMultipleErrorPropagatesSelectedReceiverMessage() {
+ final TSStatus status =
+ IoTDBDataNodeReceiver.STATEMENT_STATUS_VISITOR.process(
+ new InsertRowsStatement(),
+ new TSStatus(TSStatusCode.MULTIPLE_ERROR.getStatusCode())
+ .setSubStatus(
+ Arrays.asList(
+ StatusUtils.OK,
+ new
TSStatus(TSStatusCode.WRITE_PROCESS_REJECT.getStatusCode())
+ .setMessage("receiver write queue is full"))));
+
+ Assert.assertEquals(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+ status.getCode());
+ Assert.assertEquals("receiver write queue is full", status.getMessage());
+ }
+
@Test
public void testLoadTemporaryUnavailableClassification() throws Exception {
final File tsFile = File.createTempFile("temporary-unavailable",
".tsfile");
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSinkTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSinkTest.java
new file mode 100644
index 00000000000..cd0213bc735
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSinkTest.java
@@ -0,0 +1,66 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.pipe.sink.protocol.thrift.async;
+
+import org.apache.iotdb.pipe.api.event.Event;
+import org.apache.iotdb.pipe.api.exception.PipeException;
+
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+public class IoTDBDataRegionAsyncSinkTest {
+
+ @Test
+ public void testRetryQueueFailureMessageIncludesRootCauseAndIsCleared() {
+ final IoTDBDataRegionAsyncSink sink = new IoTDBDataRegionAsyncSink();
+ final Event event = Mockito.mock(Event.class);
+
+ sink.addFailureEventToRetryQueue(
+ event,
+ new PipeException(
+ "sink transfer wrapper", new IllegalStateException("receiver
rejected request")));
+
+ Assert.assertEquals("receiver rejected request",
sink.getLastRetryFailureMessage());
+ Assert.assertTrue(
+ IoTDBDataRegionAsyncSink.formatRetryQueueFailureMessage(
+ 1, 1, 0, sink.getLastRetryFailureMessage())
+ .contains("receiver rejected request"));
+
+ sink.clearRetryEventsReferenceCount();
+
+ Assert.assertNull(sink.getLastRetryFailureMessage());
+ Assert.assertFalse(
+ IoTDBDataRegionAsyncSink.formatRetryQueueFailureMessage(
+ 0, 0, 0, sink.getLastRetryFailureMessage())
+ .contains("receiver rejected request"));
+ }
+
+ @Test
+ public void
testRetryQueueFailureMessageKeepsRootCauseTypeWhenMessageIsMissing() {
+ final IoTDBDataRegionAsyncSink sink = new IoTDBDataRegionAsyncSink();
+ final Event event = Mockito.mock(Event.class);
+
+ sink.addFailureEventToRetryQueue(
+ event, new PipeException("sink transfer wrapper", new
NullPointerException()));
+
+ Assert.assertEquals("java.lang.NullPointerException",
sink.getLastRetryFailureMessage());
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandlerCleanupTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandlerCleanupTest.java
index aca597f0463..892295b34b3 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandlerCleanupTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandlerCleanupTest.java
@@ -19,15 +19,21 @@
package org.apache.iotdb.db.pipe.sink.protocol.thrift.async.handler;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
+import org.apache.iotdb.commons.pipe.receiver.PipeReceiverStatusHandler;
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import
org.apache.iotdb.db.pipe.sink.protocol.thrift.async.IoTDBDataRegionAsyncSink;
+import org.apache.iotdb.rpc.TSStatusCode;
+import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp;
import org.junit.Assert;
import org.junit.Test;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import java.io.File;
+import java.lang.reflect.Field;
import java.nio.file.Files;
import java.util.Collections;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -70,6 +76,56 @@ public class PipeTransferTsFileHandlerCleanupTest {
}
}
+ @Test
+ public void testSealFailurePassesNestedReceiverMessageToRetryQueue() throws
Exception {
+ final File file = Files.createTempFile("pipe-transfer-seal-failure",
".tsfile").toFile();
+ try {
+ final PipeTsFileInsertionEvent event =
Mockito.mock(PipeTsFileInsertionEvent.class);
+ final IoTDBDataRegionAsyncSink sink =
Mockito.mock(IoTDBDataRegionAsyncSink.class);
+ Mockito.when(sink.statusHandler())
+ .thenReturn(new PipeReceiverStatusHandler(false, 60, false, 60,
false, false));
+
+ final PipeTransferTsFileHandler handler =
+ new PipeTransferTsFileHandler(
+ sink,
+ Collections.emptyMap(),
+ Collections.singletonList(event),
+ new AtomicInteger(1),
+ new AtomicBoolean(false),
+ file,
+ null,
+ false,
+ null);
+ markSealSignalSent(handler);
+
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("aggregate load failure")
+ .setSubStatus(
+ Collections.singletonList(
+ new
TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode())
+ .setMessage("receiver disk is full")));
+
+ Assert.assertFalse(handler.onCompleteInternal(new
TPipeTransferResp(status)));
+
+ final ArgumentCaptor<Exception> exceptionCaptor =
ArgumentCaptor.forClass(Exception.class);
+ Mockito.verify(sink)
+ .addFailureEventsToRetryQueue(
+ Mockito.eq(Collections.singletonList(event)),
exceptionCaptor.capture());
+ Assert.assertEquals("receiver disk is full",
exceptionCaptor.getValue().getMessage());
+ } finally {
+ if (file.exists()) {
+ Assert.assertTrue(file.delete());
+ }
+ }
+ }
+
+ private static void markSealSignalSent(final PipeTransferTsFileHandler
handler) throws Exception {
+ final Field field =
PipeTransferTsFileHandler.class.getDeclaredField("isSealSignalSent");
+ field.setAccessible(true);
+ ((AtomicBoolean) field.get(handler)).set(true);
+ }
+
private PipeTransferTsFileHandler createHandler(final File file, final
EnrichedEvent event)
throws Exception {
return new PipeTransferTsFileHandler(
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/PipeReceiverStatusHandler.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/PipeReceiverStatusHandler.java
index 2cb31c81fd0..55d50e76da0 100644
---
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/PipeReceiverStatusHandler.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/PipeReceiverStatusHandler.java
@@ -40,8 +40,10 @@ import javax.annotation.Nullable;
import java.util.Arrays;
import java.util.Collections;
+import java.util.IdentityHashMap;
import java.util.List;
import java.util.Objects;
+import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
@@ -101,7 +103,8 @@ public class PipeReceiverStatusHandler {
*
* @throws PipeRuntimeSinkNonReportTimeConfigurableException to retry the
current {@link Event}
* @param status the {@link TSStatus} to judge
- * @param exceptionMessage The exception message to throw
+ * @param exceptionMessage the fallback exception message when {@code
status} does not contain a
+ * usable receiver message
* @param recordMessage The message to record an ignored {@link Event}, the
caller should assure
* that the same {@link Event} generates always the same record message,
for instance, do not
* put any time-related info here
@@ -112,10 +115,15 @@ public class PipeReceiverStatusHandler {
final String recordMessage,
final boolean log4NoPrivileges) {
+ // Batch responses may put the actual receiver error only in a nested
sub-status, while callers
+ // may supply a generic transfer wrapper. Prefer the receiver message
before constructing the
+ // retry exception so the downstream error can be reported to users.
+ final String effectiveExceptionMessage =
getEffectiveExceptionMessage(status, exceptionMessage);
+
if (RetryUtils.needRetryForWrite(status.getCode())) {
LOGGER.info(PipeMessages.IOT_CONSENSUS_RETRY_WITH_INTERVAL, status);
throw new IoTConsensusV2RetryWithIncreasingIntervalException(
- exceptionMessage, Integer.MAX_VALUE);
+ effectiveExceptionMessage, Integer.MAX_VALUE);
}
if (RetryUtils.notNeedRetryForConsensus(status.getCode())) {
@@ -125,10 +133,13 @@ public class PipeReceiverStatusHandler {
if (!PipeStopStrategy.accept(null, status)) {
PipeLogger.log(
- LOGGER::info, PipeMessages.TEMPORARY_UNAVAILABLE_RETRY, status,
exceptionMessage);
+ LOGGER::info,
+ PipeMessages.TEMPORARY_UNAVAILABLE_RETRY,
+ status,
+ effectiveExceptionMessage);
final PipeResourceFailureType failureType =
PipeStopStrategy.getResourceFailureType(null, status);
- throw new PipeRuntimeSinkResourceException(exceptionMessage,
failureType);
+ throw new PipeRuntimeSinkResourceException(effectiveExceptionMessage,
failureType);
}
switch (status.getCode()) {
@@ -183,7 +194,7 @@ public class PipeReceiverStatusHandler {
status);
exceptionEventHasBeenRetried.set(true);
throw new PipeRuntimeSinkNonReportTimeConfigurableException(
- exceptionMessage,
+ effectiveExceptionMessage,
status.getCode() == 1815
&&
PipeConfig.getInstance().isPipeRetryLocallyForParallelOrUserConflict()
? Long.MAX_VALUE
@@ -201,11 +212,12 @@ public class PipeReceiverStatusHandler {
}
return;
}
- handleOtherExceptions(status, exceptionMessage, recordMessage, true);
+ handleOtherExceptions(status, effectiveExceptionMessage,
recordMessage, true);
break;
default:
// Some auth error may be wrapped in other codes
- if (Objects.nonNull(exceptionMessage) &&
exceptionMessage.contains(NO_PERMISSION_STR)) {
+ if (Objects.nonNull(effectiveExceptionMessage)
+ && effectiveExceptionMessage.contains(NO_PERMISSION_STR)) {
if (skipIfNoPrivileges) {
if (log4NoPrivileges && LOGGER.isWarnEnabled()) {
LOGGER.warn(
@@ -216,15 +228,224 @@ public class PipeReceiverStatusHandler {
}
return;
}
- handleOtherExceptions(status, exceptionMessage, recordMessage, true);
+ handleOtherExceptions(status, effectiveExceptionMessage,
recordMessage, true);
break;
}
// Other exceptions
- handleOtherExceptions(status, exceptionMessage, recordMessage, false);
+ handleOtherExceptions(status, effectiveExceptionMessage,
recordMessage, false);
break;
}
}
+ private static String getEffectiveExceptionMessage(
+ final TSStatus status, final String exceptionMessage) {
+ final String statusMessage = getStatusMessage(status);
+ if (hasText(statusMessage)) {
+ return statusMessage;
+ }
+ return exceptionMessage;
+ }
+
+ /**
+ * Returns the most useful non-blank message in a status tree. A nested
failed status is preferred
+ * over its aggregate wrapper because it usually contains the actual
receiver error. Messages used
+ * to carry redirection device paths are never reported as errors.
+ */
+ public static String getStatusMessage(final @Nullable TSStatus status) {
+ if (status == null) {
+ return null;
+ }
+
+ final StatusMessageCandidate candidate =
+ findStatusMessage(
+ status,
+ Collections.newSetFromMap(new IdentityHashMap<>()),
+ /* inheritedClassificationCode= */ -1,
+ /* depth= */ 0,
+ new int[] {0});
+ if (candidate != null) {
+ return candidate.message;
+ }
+
+ // A successful status can carry an application-specific message. It is
safe to expose one only
+ // when the whole tree has no failure status; a success child must never
mask a failure.
+ if (containsFailureStatus(status, Collections.newSetFromMap(new
IdentityHashMap<>()))) {
+ return null;
+ }
+ return findSuccessStatusMessage(status, Collections.newSetFromMap(new
IdentityHashMap<>()));
+ }
+
+ private static boolean containsFailureStatus(
+ final TSStatus status, final Set<TSStatus> visitedStatuses) {
+ if (status == null || !visitedStatuses.add(status)) {
+ return false;
+ }
+ if (isFailureStatus(status.getCode())) {
+ return true;
+ }
+ if (status.isSetSubStatus() && status.getSubStatus() != null) {
+ for (final TSStatus subStatus : status.getSubStatus()) {
+ if (containsFailureStatus(subStatus, visitedStatuses)) {
+ return true;
+ }
+ }
+ }
+ return false;
+ }
+
+ private static String findSuccessStatusMessage(
+ final TSStatus status, final Set<TSStatus> visitedStatuses) {
+ if (status == null || !visitedStatuses.add(status)) {
+ return null;
+ }
+ if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+ && hasText(status.getMessage())
+ && !isRedirectionMessageStatus(status)) {
+ return status.getMessage();
+ }
+ if (status.isSetSubStatus() && status.getSubStatus() != null) {
+ for (final TSStatus subStatus : status.getSubStatus()) {
+ final String message = findSuccessStatusMessage(subStatus,
visitedStatuses);
+ if (hasText(message)) {
+ return message;
+ }
+ }
+ }
+ return null;
+ }
+
+ private static StatusMessageCandidate findStatusMessage(
+ final TSStatus status,
+ final Set<TSStatus> visitedStatuses,
+ final int inheritedClassificationCode,
+ final int depth,
+ final int[] traversalOrder) {
+ if (status == null || !visitedStatuses.add(status)) {
+ return null;
+ }
+
+ final int currentOrder = traversalOrder[0]++;
+ final boolean wrapper = isPipeStatusWrapper(status.getCode());
+ final int classificationCode =
+ STATUS_PRIORITY.contains(status.getCode()) ? status.getCode() :
inheritedClassificationCode;
+
+ StatusMessageCandidate bestCandidate = null;
+ boolean hasFailureDescendant = false;
+ if (status.isSetSubStatus() && status.getSubStatus() != null) {
+ for (final TSStatus subStatus : status.getSubStatus()) {
+ final StatusMessageCandidate candidate =
+ findStatusMessage(
+ subStatus, visitedStatuses, classificationCode, depth + 1,
traversalOrder);
+ if (candidate != null) {
+ hasFailureDescendant = true;
+ bestCandidate = chooseBetterStatusMessage(bestCandidate, candidate);
+ }
+ }
+ }
+
+ final String ownMessage = getOwnFailureStatusMessage(status);
+ if (ownMessage != null) {
+ bestCandidate =
+ chooseBetterStatusMessage(
+ bestCandidate,
+ new StatusMessageCandidate(
+ ownMessage,
+ wrapper,
+ !hasFailureDescendant,
+ classificationCode,
+ depth,
+ currentOrder));
+ }
+ return bestCandidate;
+ }
+
+ private static StatusMessageCandidate chooseBetterStatusMessage(
+ final @Nullable StatusMessageCandidate current, final
StatusMessageCandidate candidate) {
+ if (current == null) {
+ return candidate;
+ }
+
+ // A concrete failure leaf is the closest representation of the receiver
error. Aggregate and
+ // Pipe classification statuses are retained as a fallback for responses
that have no leaf
+ // message of their own.
+ if (candidate.leaf != current.leaf) {
+ return candidate.leaf ? candidate : current;
+ }
+ if (candidate.wrapper != current.wrapper) {
+ return candidate.wrapper ? current : candidate;
+ }
+
+ final int candidatePriority =
getStatusPriority(candidate.classificationCode);
+ final int currentPriority = getStatusPriority(current.classificationCode);
+ if (candidatePriority != currentPriority) {
+ return candidatePriority > currentPriority ? candidate : current;
+ }
+ if (candidate.depth != current.depth) {
+ return candidate.depth > current.depth ? candidate : current;
+ }
+ return candidate.traversalOrder < current.traversalOrder ? candidate :
current;
+ }
+
+ private static String getOwnFailureStatusMessage(final TSStatus status) {
+ return hasText(status.getMessage())
+ && !isRedirectionMessageStatus(status)
+ && isFailureStatus(status.getCode())
+ ? status.getMessage()
+ : null;
+ }
+
+ private static boolean isPipeStatusWrapper(final int statusCode) {
+ return statusCode == TSStatusCode.MULTIPLE_ERROR.getStatusCode()
+ || statusCode ==
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()
+ || statusCode ==
TSStatusCode.PIPE_RECEIVER_IDEMPOTENT_CONFLICT_EXCEPTION.getStatusCode()
+ || statusCode ==
TSStatusCode.PIPE_RECEIVER_USER_CONFLICT_EXCEPTION.getStatusCode()
+ || statusCode
+ ==
TSStatusCode.PIPE_RECEIVER_PARALLEL_OR_USER_CONFLICT_EXCEPTION.getStatusCode();
+ }
+
+ private static int getStatusPriority(final int statusCode) {
+ return STATUS_PRIORITY.indexOf(statusCode);
+ }
+
+ private static boolean isRedirectionMessageStatus(final TSStatus status) {
+ return status.getCode() ==
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()
+ || (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+ && status.isSetRedirectNode());
+ }
+
+ private static boolean isFailureStatus(final int statusCode) {
+ return statusCode != TSStatusCode.SUCCESS_STATUS.getStatusCode()
+ && statusCode != TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode();
+ }
+
+ private static boolean hasText(final String message) {
+ return message != null && !message.trim().isEmpty();
+ }
+
+ private static final class StatusMessageCandidate {
+ private final String message;
+ private final boolean wrapper;
+ private final boolean leaf;
+ private final int classificationCode;
+ private final int depth;
+ private final int traversalOrder;
+
+ private StatusMessageCandidate(
+ final String message,
+ final boolean wrapper,
+ final boolean leaf,
+ final int classificationCode,
+ final int depth,
+ final int traversalOrder) {
+ this.message = message;
+ this.wrapper = wrapper;
+ this.leaf = leaf;
+ this.classificationCode = classificationCode;
+ this.depth = depth;
+ this.traversalOrder = traversalOrder;
+ }
+ }
+
private synchronized void handleOtherExceptions(
final TSStatus status,
final String exceptionMessage,
@@ -354,6 +575,10 @@ public class PipeReceiverStatusHandler {
}
}
resultStatus.setSubStatus(givenStatusList);
+ final String statusMessage = getStatusMessage(resultStatus);
+ if (hasText(statusMessage)) {
+ resultStatus.setMessage(statusMessage);
+ }
return resultStatus;
}
diff --git
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/PipeReceiverStatusHandlerTest.java
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/PipeReceiverStatusHandlerTest.java
index 38d2e015165..6a0792b707a 100644
---
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/PipeReceiverStatusHandlerTest.java
+++
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/PipeReceiverStatusHandlerTest.java
@@ -19,7 +19,9 @@
package org.apache.iotdb.commons.pipe.datastructure;
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import
org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkNonReportTimeConfigurableException;
import org.apache.iotdb.commons.pipe.receiver.PipeReceiverStatusHandler;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -29,6 +31,9 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.Marker;
+import java.util.Arrays;
+import java.util.Collections;
+
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_CONFLICT_RECORD_IGNORED_DATA_DEFAULT_VALUE;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_CONFLICT_RESOLVE_STRATEGY_DEFAULT_VALUE;
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_CONFLICT_RETRY_MAX_TIME_SECONDS_DEFAULT_VALUE;
@@ -36,6 +41,227 @@ import static
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CON
import static
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_OTHERS_RETRY_MAX_TIME_SECONDS_DEFAULT_VALUE;
public class PipeReceiverStatusHandlerTest {
+ @Test
+ public void testNestedStatusMessageOverridesGenericExceptionMessage() {
+ final PipeReceiverStatusHandler handler =
+ new PipeReceiverStatusHandler(false, 60, false, 60, false, false);
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("aggregate receiver error")
+ .setSubStatus(
+ Arrays.asList(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()),
+ new TSStatus(TSStatusCode.METADATA_ERROR.getStatusCode())
+ .setMessage("receiver rejected request")));
+
+ try {
+ handler.handle(status, "generic sink transfer error", "record");
+ Assert.fail("Expected a retry exception");
+ } catch (final PipeRuntimeSinkNonReportTimeConfigurableException e) {
+ Assert.assertEquals("receiver rejected request", e.getMessage());
+ }
+ }
+
+ @Test
+ public void testStatusMessagePrefersNestedFailureOverOuterAggregateMessage()
{
+ final TSStatus status =
+ new TSStatus(TSStatusCode.MULTIPLE_ERROR.getStatusCode())
+ .setMessage("outer message")
+ .setSubStatus(
+ Arrays.asList(
+ new TSStatus(TSStatusCode.METADATA_ERROR.getStatusCode())
+ .setMessage("inner message")));
+
+ Assert.assertEquals("inner message",
PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
+ @Test
+ public void testStatusMessagePrefersConcreteFailureOverClassifiedWrapper() {
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("outer transfer failure")
+ .setSubStatus(
+ Arrays.asList(
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION
+ .getStatusCode())
+ .setMessage("generic receiver failure"),
+ new TSStatus(TSStatusCode.METADATA_ERROR.getStatusCode())
+ .setMessage("actual metadata failure")));
+
+ Assert.assertEquals(
+ "actual metadata failure",
PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
+ @Test
+ public void testStatusMessageFindsDeepNestedFailure() {
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("outer transfer failure")
+ .setSubStatus(
+ Collections.singletonList(
+ new TSStatus(TSStatusCode.MULTIPLE_ERROR.getStatusCode())
+ .setMessage("batch failure")
+ .setSubStatus(
+ Collections.singletonList(
+ new
TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode())
+ .setMessage("receiver disk is full")))));
+
+ Assert.assertEquals(
+ "receiver disk is full",
PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
+ @Test
+ public void testStatusMessageUsesHighestPriorityClassifiedFailure() {
+ final TSStatus status =
+ new TSStatus(TSStatusCode.MULTIPLE_ERROR.getStatusCode())
+ .setSubStatus(
+ Arrays.asList(
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_IDEMPOTENT_CONFLICT_EXCEPTION
+ .getStatusCode())
+ .setMessage("idempotent conflict"),
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_USER_CONFLICT_EXCEPTION.getStatusCode())
+ .setMessage("user conflict"),
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION
+ .getStatusCode())
+ .setMessage("receiver unavailable")));
+
+ Assert.assertEquals("receiver unavailable",
PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
+ @Test
+ public void testStatusMessagePrefersFailureOverRedirectionMessage() {
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setSubStatus(
+ Arrays.asList(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())
+ .setRedirectNode(new TEndPoint("127.0.0.1", 6667))
+ .setMessage("root.sg.device"),
+ new TSStatus(TSStatusCode.METADATA_ERROR.getStatusCode())
+ .setMessage("receiver rejected request")));
+
+ Assert.assertEquals(
+ "receiver rejected request",
PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
+ @Test
+ public void testStatusMessageDoesNotUseRedirectionDevicePath() {
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setSubStatus(
+ Arrays.asList(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())
+ .setRedirectNode(new TEndPoint("127.0.0.1", 6667))
+ .setMessage("root.sg.device")));
+
+ Assert.assertNull(PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
+ @Test
+ public void testStatusMessageIgnoresSuccessAndRedirectionMessages() {
+ Assert.assertNull(
+ PipeReceiverStatusHandler.getStatusMessage(
+ new TSStatus(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode())
+ .setMessage("root.sg.device")));
+
+ final TSStatus status =
+ new
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setSubStatus(
+ Arrays.asList(
+ new
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()).setMessage("ok"),
+ new
TSStatus(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode())
+ .setMessage("root.sg.redirected")));
+ Assert.assertNull(PipeReceiverStatusHandler.getStatusMessage(status));
+
+ final PipeReceiverStatusHandler handler =
+ new PipeReceiverStatusHandler(false, 60, false, 60, false, false);
+ try {
+ handler.handle(status, "generic sink transfer error", "record");
+ Assert.fail("Expected a retry exception");
+ } catch (final PipeRuntimeSinkNonReportTimeConfigurableException e) {
+ Assert.assertEquals("generic sink transfer error", e.getMessage());
+ }
+ }
+
+ @Test
+ public void testStatusMessageKeepsCustomMessageOnSuccessfulRoot() {
+ Assert.assertEquals(
+ "receiver accepted the request",
+ PipeReceiverStatusHandler.getStatusMessage(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())
+ .setMessage("receiver accepted the request")));
+ }
+
+ @Test
+ public void testPriorStatusUsesLaterMessageForSameCode() {
+ final TSStatus status =
+ PipeReceiverStatusHandler.getPriorStatus(
+ Arrays.asList(
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()),
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("receiver disk is full")));
+
+ Assert.assertEquals(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+ status.getCode());
+ Assert.assertEquals("receiver disk is full", status.getMessage());
+ }
+
+ @Test
+ public void testPriorStatusPreservesCustomSuccessMessage() {
+ final TSStatus status =
+ PipeReceiverStatusHandler.getPriorStatus(
+ Collections.singletonList(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())
+ .setMessage("receiver accepted the request")));
+
+ Assert.assertEquals("receiver accepted the request", status.getMessage());
+ }
+
+ @Test
+ public void testPriorStatusPrefersNestedConcreteFailureMessage() {
+ final TSStatus status =
+ PipeReceiverStatusHandler.getPriorStatus(
+ Arrays.asList(
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("generic receiver failure")
+ .setSubStatus(
+ Collections.singletonList(
+ new
TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode())
+ .setMessage("receiver disk is full"))),
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())));
+
+ Assert.assertEquals("receiver disk is full", status.getMessage());
+ }
+
+ @Test
+ public void testPriorStatusPropagatesSelectedFailureMessage() {
+ final TSStatus status =
+ PipeReceiverStatusHandler.getPriorStatus(
+ Arrays.asList(
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_IDEMPOTENT_CONFLICT_EXCEPTION.getStatusCode())
+ .setMessage("receiver already contains this point"),
+ new
TSStatus(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode())
+ .setMessage("root.sg.device"),
+ new TSStatus(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+ .setMessage("receiver disk is full")));
+
+ Assert.assertEquals(
+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+ status.getCode());
+ Assert.assertEquals("receiver disk is full", status.getMessage());
+ Assert.assertEquals(
+ "receiver disk is full",
PipeReceiverStatusHandler.getStatusMessage(status));
+ }
+
@Test
public void testAuthLogger() {
final PipeReceiverStatusHandler handler =