Skip to content

Commit 18ec5f0

Browse files
authored
[Pipe] Preserve downstream sink error messages (#18553)
1 parent 6b0354f commit 18ec5f0

8 files changed

Lines changed: 681 additions & 16 deletions

File tree

‎iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2590,6 +2590,13 @@ private DataNodePipeMessages() {}
25902590
public static final String MESSAGE_TRANSFER_FILE_ARG_ERROR_RESULT_STATUS_ARG_E565D9FD =
25912591
"Transfer file %s error, result status %s.";
25922592

2593+
public static final String
2594+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_5B4B2E7C =
2595+
"Failed to retry transferring events in the retry queue. Remaining events: %d (tablet events: %d, tsfile events: %d).";
2596+
public static final String
2597+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_LAST_FAILURE_ARG_EB8F9DCD =
2598+
"Failed to retry transferring events in the retry queue. Remaining events: %d (tablet events: %d, tsfile events: %d). Last failure: %s.";
2599+
25932600
public static final String EXCEPTION_LEGACY_PIPE_RECEIVER_REQUIRES_A_LOGGED_IN_SESSION_D96219BF =
25942601
"Legacy pipe receiver requires a logged-in session.";
25952602
public static final String EXCEPTION_FAILED_TO_SET_UP_CONSENSUS_SUBSCRIPTION_FOR_TOPIC_ARG_IN_CONSUMER_GROUP_ARG_ARG_A7FA88F3 =

‎iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2416,6 +2416,13 @@ private DataNodePipeMessages() {}
24162416
public static final String MESSAGE_TRANSFER_FILE_ARG_ERROR_RESULT_STATUS_ARG_E565D9FD =
24172417
"传输文件 %s 出错,结果状态为 %s。";
24182418

2419+
public static final String
2420+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_5B4B2E7C =
2421+
"重试 retry queue 中的事件失败。剩余事件:%d(tablet 事件:%d,tsfile 事件:%d)。";
2422+
public static final String
2423+
EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_LAST_FAILURE_ARG_EB8F9DCD =
2424+
"重试 retry queue 中的事件失败。剩余事件:%d(tablet 事件:%d,tsfile 事件:%d)。最近一次失败:%s。";
2425+
24192426
public static final String EXCEPTION_LEGACY_PIPE_RECEIVER_REQUIRES_A_LOGGED_IN_SESSION_D96219BF =
24202427
"Legacy pipe receiver 需要已登录的 session。";
24212428
public static final String EXCEPTION_FAILED_TO_SET_UP_CONSENSUS_SUBSCRIPTION_FOR_TOPIC_ARG_IN_CONSUMER_GROUP_ARG_ARG_A7FA88F3 =

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java‎

Lines changed: 67 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
import org.apache.iotdb.commons.pipe.resource.log.PipeLogger;
3535
import org.apache.iotdb.commons.pipe.sink.protocol.IoTDBSink;
3636
import org.apache.iotdb.commons.pipe.sink.protocol.PipeSinkWithSchedulingDelay;
37+
import org.apache.iotdb.commons.utils.ErrorHandlingCommonUtils;
3738
import org.apache.iotdb.db.i18n.DataNodePipeMessages;
3839
import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
3940
import org.apache.iotdb.db.pipe.event.common.deletion.PipeDeleteDataNodeEvent;
@@ -138,6 +139,8 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink implements PipeSinkWithS
138139
// Guarded by this. Events need identity semantics because the same payload may compare equal.
139140
private final Map<Event, PipeResourceFailureType> retryEvent2ResourceFailureType =
140141
new IdentityHashMap<>();
142+
// Keep only the latest text to avoid retaining the complete exception chain for every event.
143+
private volatile String lastRetryFailureMessage;
141144

142145
private IoTDBDataNodeAsyncClientManager clientManager;
143146
private IoTDBDataNodeAsyncClientManager transferTsFileClientManager;
@@ -732,13 +735,11 @@ private void transferQueuedEventsIfNecessary(final boolean forced) {
732735

733736
if (remainingEvents <= retryEventQueue.size() + retryTsFileQueue.size()) {
734737
final String message =
735-
"Failed to retry transferring events in the retry queue. Remaining events: "
736-
+ (retryEventQueue.size() + retryTsFileQueue.size())
737-
+ " (tablet events: "
738-
+ retryEventQueueEventCounter.getTabletInsertionEventCount()
739-
+ ", tsfile events: "
740-
+ retryEventQueueEventCounter.getTsFileInsertionEventCount()
741-
+ ").";
738+
formatRetryQueueFailureMessage(
739+
retryEventQueue.size() + retryTsFileQueue.size(),
740+
retryEventQueueEventCounter.getTabletInsertionEventCount(),
741+
retryEventQueueEventCounter.getTsFileInsertionEventCount(),
742+
lastRetryFailureMessage);
742743
final PipeResourceFailureType retryQueueResourceFailureType =
743744
getRetryQueueResourceFailureType();
744745
if (retryQueueResourceFailureType != null) {
@@ -751,6 +752,12 @@ private void transferQueuedEventsIfNecessary(final boolean forced) {
751752
}
752753
}
753754
}
755+
756+
synchronized (this) {
757+
if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
758+
lastRetryFailureMessage = null;
759+
}
760+
}
754761
}
755762

756763
private void retryTransfer(final TabletInsertionEvent tabletInsertionEvent) {
@@ -834,6 +841,14 @@ private synchronized void addFailureEventToRetryQueue(
834841
return;
835842
}
836843

844+
if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
845+
lastRetryFailureMessage = null;
846+
}
847+
848+
if (e != null) {
849+
lastRetryFailureMessage = getRetryFailureMessage(e);
850+
}
851+
837852
if (resourceFailureType != null && event instanceof EnrichedEvent) {
838853
final EnrichedEvent enrichedEvent = (EnrichedEvent) event;
839854
final Pair<String, Long> pipeKey =
@@ -881,6 +896,46 @@ public void addFailureEventsToRetryQueue(
881896
events.forEach(event -> addFailureEventToRetryQueue(event, e, failureRecordedPipes));
882897
}
883898

899+
static String formatRetryQueueFailureMessage(
900+
final int remainingEvents,
901+
final int tabletEventCount,
902+
final int tsFileEventCount,
903+
final String lastFailureMessage) {
904+
if (!hasText(lastFailureMessage)) {
905+
return String.format(
906+
DataNodePipeMessages
907+
.EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_5B4B2E7C,
908+
remainingEvents,
909+
tabletEventCount,
910+
tsFileEventCount);
911+
}
912+
return String.format(
913+
DataNodePipeMessages
914+
.EXCEPTION_FAILED_TO_RETRY_TRANSFERRING_EVENTS_IN_THE_RETRY_QUEUE_REMAINING_EVENTS_ARG_TABLET_EVENTS_ARG_TSFILE_EVENTS_ARG_LAST_FAILURE_ARG_EB8F9DCD,
915+
remainingEvents,
916+
tabletEventCount,
917+
tsFileEventCount,
918+
lastFailureMessage);
919+
}
920+
921+
private static String getRetryFailureMessage(final Exception exception) {
922+
final Throwable rootCause = ErrorHandlingCommonUtils.getRootCause(exception);
923+
if (hasText(rootCause.getMessage())) {
924+
return rootCause.getMessage();
925+
}
926+
// Throwable#getMessage() is null for exceptions such as a bare NPE. Keep the type in the
927+
// reported sink error instead of falling back to a generic transfer wrapper.
928+
return rootCause.toString();
929+
}
930+
931+
private static boolean hasText(final String message) {
932+
return message != null && !message.trim().isEmpty();
933+
}
934+
935+
synchronized String getLastRetryFailureMessage() {
936+
return lastRetryFailureMessage;
937+
}
938+
884939
private synchronized PipeResourceFailureType getRetryQueueResourceFailureType() {
885940
for (final PipeResourceFailureType failureType : PipeResourceFailureType.values()) {
886941
if (retryEvent2ResourceFailureType.containsValue(failureType)) {
@@ -1056,6 +1111,10 @@ && isDroppedPipe((EnrichedEvent) event, committerKey)) {
10561111
}
10571112
return false;
10581113
});
1114+
1115+
if (retryEventQueue.isEmpty() && retryTsFileQueue.isEmpty()) {
1116+
lastRetryFailureMessage = null;
1117+
}
10591118
}
10601119

10611120
@Override
@@ -1110,6 +1169,7 @@ public synchronized void clearRetryEventsReferenceCount() {
11101169
}
11111170
}
11121171
retryEvent2ResourceFailureType.clear();
1172+
lastRetryFailureMessage = null;
11131173
}
11141174

11151175
//////////////////////// APIs provided for metric framework ////////////////////////

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/PipeStatementTsStatusVisitorTest.java‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,24 @@ StatusUtils.OK, new TSStatus(TSStatusCode.OUT_OF_TTL.getStatusCode()))))
6666
.getCode());
6767
}
6868

69+
@Test
70+
public void testMultipleErrorPropagatesSelectedReceiverMessage() {
71+
final TSStatus status =
72+
IoTDBDataNodeReceiver.STATEMENT_STATUS_VISITOR.process(
73+
new InsertRowsStatement(),
74+
new TSStatus(TSStatusCode.MULTIPLE_ERROR.getStatusCode())
75+
.setSubStatus(
76+
Arrays.asList(
77+
StatusUtils.OK,
78+
new TSStatus(TSStatusCode.WRITE_PROCESS_REJECT.getStatusCode())
79+
.setMessage("receiver write queue is full"))));
80+
81+
Assert.assertEquals(
82+
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
83+
status.getCode());
84+
Assert.assertEquals("receiver write queue is full", status.getMessage());
85+
}
86+
6987
@Test
7088
public void testLoadTemporaryUnavailableClassification() throws Exception {
7189
final File tsFile = File.createTempFile("temporary-unavailable", ".tsfile");
Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,66 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.db.pipe.sink.protocol.thrift.async;
21+
22+
import org.apache.iotdb.pipe.api.event.Event;
23+
import org.apache.iotdb.pipe.api.exception.PipeException;
24+
25+
import org.junit.Assert;
26+
import org.junit.Test;
27+
import org.mockito.Mockito;
28+
29+
public class IoTDBDataRegionAsyncSinkTest {
30+
31+
@Test
32+
public void testRetryQueueFailureMessageIncludesRootCauseAndIsCleared() {
33+
final IoTDBDataRegionAsyncSink sink = new IoTDBDataRegionAsyncSink();
34+
final Event event = Mockito.mock(Event.class);
35+
36+
sink.addFailureEventToRetryQueue(
37+
event,
38+
new PipeException(
39+
"sink transfer wrapper", new IllegalStateException("receiver rejected request")));
40+
41+
Assert.assertEquals("receiver rejected request", sink.getLastRetryFailureMessage());
42+
Assert.assertTrue(
43+
IoTDBDataRegionAsyncSink.formatRetryQueueFailureMessage(
44+
1, 1, 0, sink.getLastRetryFailureMessage())
45+
.contains("receiver rejected request"));
46+
47+
sink.clearRetryEventsReferenceCount();
48+
49+
Assert.assertNull(sink.getLastRetryFailureMessage());
50+
Assert.assertFalse(
51+
IoTDBDataRegionAsyncSink.formatRetryQueueFailureMessage(
52+
0, 0, 0, sink.getLastRetryFailureMessage())
53+
.contains("receiver rejected request"));
54+
}
55+
56+
@Test
57+
public void testRetryQueueFailureMessageKeepsRootCauseTypeWhenMessageIsMissing() {
58+
final IoTDBDataRegionAsyncSink sink = new IoTDBDataRegionAsyncSink();
59+
final Event event = Mockito.mock(Event.class);
60+
61+
sink.addFailureEventToRetryQueue(
62+
event, new PipeException("sink transfer wrapper", new NullPointerException()));
63+
64+
Assert.assertEquals("java.lang.NullPointerException", sink.getLastRetryFailureMessage());
65+
}
66+
}

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandlerCleanupTest.java‎

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,15 +19,21 @@
1919

2020
package org.apache.iotdb.db.pipe.sink.protocol.thrift.async.handler;
2121

22+
import org.apache.iotdb.common.rpc.thrift.TSStatus;
2223
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
24+
import org.apache.iotdb.commons.pipe.receiver.PipeReceiverStatusHandler;
2325
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
2426
import org.apache.iotdb.db.pipe.sink.protocol.thrift.async.IoTDBDataRegionAsyncSink;
27+
import org.apache.iotdb.rpc.TSStatusCode;
28+
import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp;
2529

2630
import org.junit.Assert;
2731
import org.junit.Test;
32+
import org.mockito.ArgumentCaptor;
2833
import org.mockito.Mockito;
2934

3035
import java.io.File;
36+
import java.lang.reflect.Field;
3137
import java.nio.file.Files;
3238
import java.util.Collections;
3339
import java.util.concurrent.atomic.AtomicBoolean;
@@ -70,6 +76,56 @@ public void testCloseKeepsSourceTsFile() throws Exception {
7076
}
7177
}
7278

79+
@Test
80+
public void testSealFailurePassesNestedReceiverMessageToRetryQueue() throws Exception {
81+
final File file = Files.createTempFile("pipe-transfer-seal-failure", ".tsfile").toFile();
82+
try {
83+
final PipeTsFileInsertionEvent event = Mockito.mock(PipeTsFileInsertionEvent.class);
84+
final IoTDBDataRegionAsyncSink sink = Mockito.mock(IoTDBDataRegionAsyncSink.class);
85+
Mockito.when(sink.statusHandler())
86+
.thenReturn(new PipeReceiverStatusHandler(false, 60, false, 60, false, false));
87+
88+
final PipeTransferTsFileHandler handler =
89+
new PipeTransferTsFileHandler(
90+
sink,
91+
Collections.emptyMap(),
92+
Collections.singletonList(event),
93+
new AtomicInteger(1),
94+
new AtomicBoolean(false),
95+
file,
96+
null,
97+
false,
98+
null);
99+
markSealSignalSent(handler);
100+
101+
final TSStatus status =
102+
new TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
103+
.setMessage("aggregate load failure")
104+
.setSubStatus(
105+
Collections.singletonList(
106+
new TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode())
107+
.setMessage("receiver disk is full")));
108+
109+
Assert.assertFalse(handler.onCompleteInternal(new TPipeTransferResp(status)));
110+
111+
final ArgumentCaptor<Exception> exceptionCaptor = ArgumentCaptor.forClass(Exception.class);
112+
Mockito.verify(sink)
113+
.addFailureEventsToRetryQueue(
114+
Mockito.eq(Collections.singletonList(event)), exceptionCaptor.capture());
115+
Assert.assertEquals("receiver disk is full", exceptionCaptor.getValue().getMessage());
116+
} finally {
117+
if (file.exists()) {
118+
Assert.assertTrue(file.delete());
119+
}
120+
}
121+
}
122+
123+
private static void markSealSignalSent(final PipeTransferTsFileHandler handler) throws Exception {
124+
final Field field = PipeTransferTsFileHandler.class.getDeclaredField("isSealSignalSent");
125+
field.setAccessible(true);
126+
((AtomicBoolean) field.get(handler)).set(true);
127+
}
128+
73129
private PipeTransferTsFileHandler createHandler(final File file, final EnrichedEvent event)
74130
throws Exception {
75131
return new PipeTransferTsFileHandler(

0 commit comments

Comments
 (0)