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 6377522ce2c Fix pipe leader cache updates for multi-device redirects
(#18689)
6377522ce2c is described below
commit 6377522ce2cc70ea2244222c8886066504d207c8
Author: Caideyipi <[email protected]>
AuthorDate: Tue Sep 22 17:36:00 2026 +0800
Fix pipe leader cache updates for multi-device redirects (#18689)
---
.../protocol/thrift/IoTDBDataNodeReceiver.java | 13 ++-
.../PipeTransferTabletInsertNodeEventHandler.java | 13 ++-
.../PipeTransferTabletInsertionEventHandler.java | 6 +-
.../thrift/sync/IoTDBDataRegionSyncSink.java | 4 +
.../db/pipe/sink/util/cacher/LeaderCacheUtils.java | 47 +++++----
...peTransferTabletInsertNodeEventHandlerTest.java | 105 +++++++++++++++++++++
.../sink/util/cacher/LeaderCacheUtilsTest.java | 47 +++++++++
7 files changed, 210 insertions(+), 25 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
index f92d798cc41..8f2e8109823 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
@@ -529,7 +529,7 @@ public class IoTDBDataNodeReceiver extends
IoTDBFileReceiver {
return new TPipeTransferResp(
statement.isEmpty()
? RpcUtils.SUCCESS_STATUS
- : executeStatementAndClassifyExceptions(statement));
+ : executeStatementAndAddRedirectInfo(statement));
}
private TPipeTransferResp handleTransferTabletBinary(final
PipeTransferTabletBinaryReq req) {
@@ -537,7 +537,7 @@ public class IoTDBDataNodeReceiver extends
IoTDBFileReceiver {
return new TPipeTransferResp(
statement.isEmpty()
? RpcUtils.SUCCESS_STATUS
- : executeStatementAndClassifyExceptions(statement));
+ : executeStatementAndAddRedirectInfo(statement));
}
private TPipeTransferResp handleTransferTabletRaw(final
PipeTransferTabletRawReq req) {
@@ -1285,8 +1285,13 @@ public class IoTDBDataNodeReceiver extends
IoTDBFileReceiver {
* device path to each redirected sub-status.
*/
private TSStatus executeBatchStatementAndAddRedirectInfo(final
InsertBaseStatement statement) {
- final TSStatus result = executeStatementAndClassifyExceptions(statement,
5);
- return addRedirectInfoForBatch(statement, result, receiverId.get());
+ return addRedirectInfoForBatch(
+ statement, executeStatementAndClassifyExceptions(statement, 5),
receiverId.get());
+ }
+
+ private TSStatus executeStatementAndAddRedirectInfo(final
InsertBaseStatement statement) {
+ return addRedirectInfoForBatch(
+ statement, executeStatementAndClassifyExceptions(statement),
receiverId.get());
}
static TSStatus addRedirectInfoForBatch(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandler.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandler.java
index 56d1ce41b02..dc946c11fa2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandler.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandler.java
@@ -19,13 +19,16 @@
package org.apache.iotdb.db.pipe.sink.protocol.thrift.async.handler;
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import
org.apache.iotdb.commons.client.async.AsyncPipeDataTransferServiceClient;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
import
org.apache.iotdb.db.pipe.sink.protocol.thrift.async.IoTDBDataRegionAsyncSink;
+import org.apache.iotdb.db.pipe.sink.util.cacher.LeaderCacheUtils;
import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
import org.apache.thrift.TException;
+import org.apache.tsfile.utils.Pair;
public class PipeTransferTabletInsertNodeEventHandler
extends PipeTransferTabletInsertionEventHandler {
@@ -46,7 +49,13 @@ public class PipeTransferTabletInsertNodeEventHandler
@Override
protected void updateLeaderCache(final TSStatus status) {
- sink.updateLeaderCache(
- ((PipeInsertNodeTabletInsertionEvent) event).getDeviceId(),
status.getRedirectNode());
+ if (status.isSetRedirectNode()) {
+ sink.updateLeaderCache(
+ ((PipeInsertNodeTabletInsertionEvent) event).getDeviceId(),
status.getRedirectNode());
+ }
+ for (final Pair<String, TEndPoint> redirectPair :
+ LeaderCacheUtils.parseRecommendedRedirections(status)) {
+ sink.updateLeaderCache(redirectPair.getLeft(), redirectPair.getRight());
+ }
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertionEventHandler.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertionEventHandler.java
index 445e014ceec..85795ae1315 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertionEventHandler.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertionEventHandler.java
@@ -76,7 +76,11 @@ public abstract class
PipeTransferTabletInsertionEventHandler extends PipeTransf
.handle(response.getStatus(), response.getStatus().getMessage(),
event.toString());
}
event.decreaseReferenceCount(PipeTransferTabletInsertionEventHandler.class.getName(),
true);
- if (status.isSetRedirectNode()) {
+ // A multi-device InsertRowsNode response stores redirect endpoints in
per-device
+ // sub-statuses instead of on the top-level status.
+ if (status.isSetRedirectNode()
+ || (status.getCode() ==
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()
+ && status.isSetSubStatus())) {
updateLeaderCache(status);
}
} catch (final Exception e) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
index 49d9df02ab9..d7538907dd6 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
@@ -453,6 +453,10 @@ public class IoTDBDataRegionSyncSink extends
IoTDBDataNodeSyncSink {
// pipeInsertNodeTabletInsertionEvent.getDeviceId() is null for
InsertRowsNode
pipeInsertNodeTabletInsertionEvent.getDeviceId(),
status.getRedirectNode());
}
+ for (final Pair<String, TEndPoint> redirectPair :
+ LeaderCacheUtils.parseRecommendedRedirections(status)) {
+ clientManager.updateLeaderCache(redirectPair.getLeft(),
redirectPair.getRight());
+ }
}
private void doTransferWrapper(final PipeRawTabletInsertionEvent
pipeRawTabletInsertionEvent)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
index bb0357d3b90..7e4e026c9e2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
@@ -40,35 +40,46 @@ public class LeaderCacheUtils {
* @param status is the returned status after transferring a batch event.
* @return a list of pairs, each pair contains a device path and its
redirect endpoint.
*/
- public static List<Pair<String, TEndPoint>>
parseRecommendedRedirections(TSStatus status) {
+ public static List<Pair<String, TEndPoint>>
parseRecommendedRedirections(final TSStatus status) {
// Each top-level sub-status corresponds to one statement constructed by
the receiver. V2 batch
- // requests may contain any number of statements because rows are grouped
by database and table.
+ // requests may contain any number of statements because rows are grouped
by database and
+ // table. A direct InsertRowsNode request may instead put the per-device
redirect statuses
+ // directly at the top level.
final List<Pair<String, TEndPoint>> redirectList = new ArrayList<>();
- if (status == null || !status.isSetSubStatus()) {
+ if (status == null || status.getCode() !=
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()) {
return redirectList;
}
- for (final TSStatus subStatus : status.getSubStatus()) {
- if (subStatus == null
- || subStatus.getCode() !=
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()
- || !subStatus.isSetSubStatus()) {
- continue;
+ if (status.isSetSubStatus()) {
+ for (final TSStatus subStatus : status.getSubStatus()) {
+ if (subStatus != null) {
+ collectRedirects(subStatus, redirectList);
+ }
}
+ }
+
+ return redirectList;
+ }
- for (final TSStatus innerSubStatus : subStatus.getSubStatus()) {
- if (innerSubStatus != null
- && innerSubStatus.isSetRedirectNode()
- && innerSubStatus.isSetMessage()
- && !innerSubStatus.getMessage().isEmpty()) {
- // The receiver sets the message to a device path only when it can
safely associate the
- // redirection with a single tree-model device.
- redirectList.add(
- new Pair<>(innerSubStatus.getMessage(),
innerSubStatus.getRedirectNode()));
+ private static void collectRedirects(
+ final TSStatus status, final List<Pair<String, TEndPoint>> redirectList)
{
+ addRedirectIfPresent(redirectList, status);
+ if (status.isSetSubStatus()) {
+ for (final TSStatus subStatus : status.getSubStatus()) {
+ if (subStatus != null) {
+ collectRedirects(subStatus, redirectList);
}
}
}
+ }
- return redirectList;
+ private static void addRedirectIfPresent(
+ final List<Pair<String, TEndPoint>> redirectList, final TSStatus status)
{
+ if (status.isSetRedirectNode() && status.isSetMessage() &&
!status.getMessage().isEmpty()) {
+ // The receiver sets the message to a device path only when it can
safely associate the
+ // redirection with a single tree-model device.
+ redirectList.add(new Pair<>(status.getMessage(),
status.getRedirectNode()));
+ }
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandlerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandlerTest.java
new file mode 100644
index 00000000000..4c01ec68c16
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandlerTest.java
@@ -0,0 +1,105 @@
+/*
+ * 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.handler;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import
org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
+import
org.apache.iotdb.db.pipe.sink.protocol.thrift.async.IoTDBDataRegionAsyncSink;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
+import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp;
+
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.Arrays;
+
+public class PipeTransferTabletInsertNodeEventHandlerTest {
+
+ @Test
+ public void testUpdateLeaderCacheFromMultiDeviceRedirectStatus() {
+ final PipeInsertNodeTabletInsertionEvent event =
+ Mockito.mock(PipeInsertNodeTabletInsertionEvent.class);
+ Mockito.when(event.getDeviceId()).thenReturn(null);
+ final IoTDBDataRegionAsyncSink sink =
Mockito.mock(IoTDBDataRegionAsyncSink.class);
+ final PipeTransferTabletInsertNodeEventHandler handler =
+ new PipeTransferTabletInsertNodeEventHandler(event, null, sink);
+
+ final TEndPoint firstEndPoint = new TEndPoint("127.0.0.2", 6667);
+ final TEndPoint secondEndPoint = new TEndPoint("127.0.0.3", 6667);
+ handler.updateLeaderCache(
+ RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+ .setSubStatus(
+ Arrays.asList(
+ redirectStatus("root.sg.device1", firstEndPoint),
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS),
+ redirectStatus("root.sg.device2", secondEndPoint))));
+
+ Mockito.verify(sink).updateLeaderCache("root.sg.device1", firstEndPoint);
+ Mockito.verify(sink).updateLeaderCache("root.sg.device2", secondEndPoint);
+ Mockito.verifyNoMoreInteractions(sink);
+ }
+
+ @Test
+ public void testUpdateLeaderCacheFromSingleDeviceRedirectStatus() {
+ final PipeInsertNodeTabletInsertionEvent event =
+ Mockito.mock(PipeInsertNodeTabletInsertionEvent.class);
+ Mockito.when(event.getDeviceId()).thenReturn("root.sg.device");
+ final IoTDBDataRegionAsyncSink sink =
Mockito.mock(IoTDBDataRegionAsyncSink.class);
+ final PipeTransferTabletInsertNodeEventHandler handler =
+ new PipeTransferTabletInsertNodeEventHandler(event, null, sink);
+ final TEndPoint redirectEndPoint = new TEndPoint("127.0.0.4", 6667);
+
+ handler.updateLeaderCache(
+
RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS).setRedirectNode(redirectEndPoint));
+
+ Mockito.verify(sink).updateLeaderCache("root.sg.device", redirectEndPoint);
+ }
+
+ @Test
+ public void testOnCompleteUpdatesMultiDeviceLeaderCache() {
+ final PipeInsertNodeTabletInsertionEvent event =
+ Mockito.mock(PipeInsertNodeTabletInsertionEvent.class);
+ Mockito.when(event.getDeviceId()).thenReturn(null);
+ final IoTDBDataRegionAsyncSink sink =
Mockito.mock(IoTDBDataRegionAsyncSink.class);
+ final PipeTransferTabletInsertNodeEventHandler handler =
+ new PipeTransferTabletInsertNodeEventHandler(event, null, sink);
+ final TEndPoint redirectEndPoint = new TEndPoint("127.0.0.5", 6667);
+
+ handler.onCompleteInternal(
+ new TPipeTransferResp(
+ RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+ .setSubStatus(
+ Arrays.asList(
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS),
+ redirectStatus("root.sg.device3",
redirectEndPoint)))));
+
+ Mockito.verify(sink).updateLeaderCache("root.sg.device3",
redirectEndPoint);
+ Mockito.verify(event)
+
.decreaseReferenceCount(PipeTransferTabletInsertionEventHandler.class.getName(),
true);
+ }
+
+ private static TSStatus redirectStatus(final String deviceId, final
TEndPoint endPoint) {
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)
+ .setMessage(deviceId)
+ .setRedirectNode(endPoint);
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
index 62d5b57cedf..cabe9015bfd 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
@@ -60,6 +60,53 @@ public class LeaderCacheUtilsTest {
Assert.assertEquals(redirectEndPoint, redirects.get(0).getRight());
}
+ @Test
+ public void testParseRecommendedRedirectionsFromDirectMultiDeviceStatus() {
+ final TEndPoint redirectEndPoint = new TEndPoint("127.0.0.3", 6667);
+ final TSStatus directStatus =
+ RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+ .setSubStatus(
+ Arrays.asList(
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS),
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)
+ .setMessage("root.sg.device2")
+ .setRedirectNode(redirectEndPoint)));
+
+ final List<Pair<String, TEndPoint>> redirects =
+ LeaderCacheUtils.parseRecommendedRedirections(directStatus);
+
+ Assert.assertEquals(
+ Collections.singletonList(new Pair<>("root.sg.device2",
redirectEndPoint)), redirects);
+ }
+
+ @Test
+ public void testParseRecommendedRedirectionsIgnoresTopLevelRedirect() {
+ final TSStatus status =
+ RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+ .setMessage("redirect recommendation")
+ .setRedirectNode(new TEndPoint("127.0.0.4", 6667));
+
+
Assert.assertTrue(LeaderCacheUtils.parseRecommendedRedirections(status).isEmpty());
+
Assert.assertTrue(LeaderCacheUtils.parseRecommendedRedirections(null).isEmpty());
+ }
+
+ @Test
+ public void testParseRecommendedRedirectionsIgnoresNonRedirectionStatus() {
+ final TSStatus status =
+ RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)
+ .setSubStatus(
+ Collections.singletonList(
+ redirectStatus("root.sg.device3", new
TEndPoint("127.0.0.5", 6667))));
+
+
Assert.assertTrue(LeaderCacheUtils.parseRecommendedRedirections(status).isEmpty());
+ }
+
+ private static TSStatus redirectStatus(final String deviceId, final
TEndPoint endPoint) {
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)
+ .setMessage(deviceId)
+ .setRedirectNode(endPoint);
+ }
+
@Test
public void testIgnoreRedirectsWithoutDevicePath() {
final TEndPoint redirectEndPoint = new TEndPoint("127.0.0.2", 6667);