This is an automated email from the ASF dual-hosted git repository. Caideyipi pushed a commit to branch fix/tdb-30-pipe-leader-cache in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 08466e298260234a3a3eac71162e713090954b97 Author: Caideyipi <[email protected]> AuthorDate: Mon Sep 21 14:29:26 2026 +0800 Fix pipe leader cache updates for multi-device redirects --- .../protocol/thrift/IoTDBDataNodeReceiver.java | 13 ++- .../PipeTransferTabletInsertNodeEventHandler.java | 13 ++- .../PipeTransferTabletInsertionEventHandler.java | 6 +- .../thrift/sync/IoTDBDataRegionSyncSink.java | 4 + .../db/pipe/sink/util/cacher/LeaderCacheUtils.java | 41 +++++--- ...peTransferTabletInsertNodeEventHandlerTest.java | 105 +++++++++++++++++++++ .../sink/util/cacher/LeaderCacheUtilsTest.java | 47 +++++++++ 7 files changed, 210 insertions(+), 19 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 895993704e2..fe3378a8973 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) { @@ -1287,7 +1287,14 @@ public class IoTDBDataNodeReceiver extends IoTDBFileReceiver { * 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 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 0f6beade80d..524a5a3a040 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 @@ 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.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<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()) { + // 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 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 76c6bef2467..161766ddca5 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 class LeaderCacheUtilsTest { 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<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); + } }
