From 08466e298260234a3a3eac71162e713090954b97 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Mon, 21 Sep 2026 14:29:26 +0800 Subject: [PATCH] Fix pipe leader cache updates for multi-device redirects --- .../thrift/IoTDBDataNodeReceiver.java | 13 ++- ...eTransferTabletInsertNodeEventHandler.java | 13 ++- ...peTransferTabletInsertionEventHandler.java | 6 +- .../thrift/sync/IoTDBDataRegionSyncSink.java | 4 + .../sink/util/cacher/LeaderCacheUtils.java | 41 ++++--- ...nsferTabletInsertNodeEventHandlerTest.java | 105 ++++++++++++++++++ .../util/cacher/LeaderCacheUtilsTest.java | 47 ++++++++ 7 files changed, 210 insertions(+), 19 deletions(-) create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTabletInsertNodeEventHandlerTest.java 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 895993704e28d..fe3378a89732d 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 @@ private TPipeTransferResp handleTransferTabletInsertNode( return new TPipeTransferResp( statement.isEmpty() ? RpcUtils.SUCCESS_STATUS - : executeStatementAndClassifyExceptions(statement)); + : executeStatementAndAddRedirectInfo(statement)); } private TPipeTransferResp handleTransferTabletBinary(final PipeTransferTabletBinaryReq req) { @@ -537,7 +537,7 @@ private TPipeTransferResp handleTransferTabletBinary(final PipeTransferTabletBin return new TPipeTransferResp( statement.isEmpty() ? RpcUtils.SUCCESS_STATUS - : executeStatementAndClassifyExceptions(statement)); + : executeStatementAndAddRedirectInfo(statement)); } private TPipeTransferResp handleTransferTabletRaw(final PipeTransferTabletRawReq req) { @@ -1287,7 +1287,14 @@ private static long parsePipeCreationTime(final String pipeCreationTime) { * message field. */ private TSStatus executeBatchStatementAndAddRedirectInfo(final InsertBaseStatement statement) { - final TSStatus result = executeStatementAndClassifyExceptions(statement, 5); + return addRedirectInfo(statement, executeStatementAndClassifyExceptions(statement, 5)); + } + + private TSStatus executeStatementAndAddRedirectInfo(final InsertBaseStatement statement) { + return addRedirectInfo(statement, executeStatementAndClassifyExceptions(statement)); + } + + private TSStatus addRedirectInfo(final InsertBaseStatement statement, final TSStatus result) { if (result.getCode() == TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode() && result.getSubStatusSize() > 0) { 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 56d1ce41b029c..dc946c11fa210 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 @@ protected void doTransfer( @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 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 445e014ceec94..85795ae1315e8 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 @@ protected boolean onCompleteInternal(final TPipeTransferResp response) { .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 49d9df02ab99d..d7538907dd693 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 @@ private void doTransfer( // pipeInsertNodeTabletInsertionEvent.getDeviceId() is null for InsertRowsNode pipeInsertNodeTabletInsertionEvent.getDeviceId(), status.getRedirectNode()); } + for (final Pair 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 0f6beade80d66..524a5a3a040f5 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,30 +40,45 @@ private 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> parseRecommendedRedirections(TSStatus status) { + public static List> 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> redirectList = new ArrayList<>(); - if (!status.isSetSubStatus()) { + if (status == null || status.getCode() != TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()) { return redirectList; } - for (final TSStatus subStatus : status.getSubStatus()) { - if (subStatus.getCode() != TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()) { - 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.isSetRedirectNode()) { - // We assume that innerSubStatus.getMessage() is a device path. - // The message field should be a device path. - redirectList.add( - new Pair<>(innerSubStatus.getMessage(), innerSubStatus.getRedirectNode())); + private static void collectRedirects( + final TSStatus status, final List> 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> redirectList, final TSStatus status) { + if (status.isSetRedirectNode() && status.isSetMessage()) { + // The receiver records the device path in the message field for redirected devices. + 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 0000000000000..4c01ec68c1633 --- /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 76c6bef24672a..161766ddca5f9 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 @@ -59,4 +59,51 @@ public void testParseRecommendedRedirectionsFromVariableStatementCount() { Assert.assertEquals("table1.device1", redirects.get(0).getLeft()); 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> 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); + } }