This is an automated email from the ASF dual-hosted git repository.

liaoxin01 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 94ad1c31bbc [fix](load) resolve adaptive random bucket routing for 
stream load (#68006)
94ad1c31bbc is described below

commit 94ad1c31bbc3b352365a75849f26147581456e0d
Author: Refrain <[email protected]>
AuthorDate: Fri Oct 9 14:51:36 2026 +0800

    [fix](load) resolve adaptive random bucket routing for stream load (#68006)
    
    ### What problem does this PR solve?
    
    Related PR: #62661
    
    Problem Summary:
    
    Cloud random distribution tables turn adaptive random bucket mode on for
    the sink
    (`enable_adaptive_random_bucket_load`, default true), and the FE is
    expected to compute, for every
    sink backend, the bucket owner backend and the receiver local bucket
    sequence through
    `OlapTableSink.computeAdaptiveRandomBucketAssignments` /
    `applyAdaptiveRandomBucketAssignments`.
    
    Only `Coordinator`, `ThriftPlansBuilder` and the runtime partition
    creation/replace paths do that.
    The stream load planning path (`NereidsStreamLoadPlanner` ->
    `StreamLoadHandler`) never did, so the
    BE received an adaptive sink without any assignment, fell back to
    treating the executing BE as the
    owner of every bucket (`vtablet_writer.cpp: adaptive_bucket_be_id`) and
    failed with
    
    ```
    unknown partition channel, load_id=..., index_id=..., partition_id=...
    ```
    
    as soon as a partition has no tablet on that BE. For a `DISTRIBUTED BY
    RANDOM BUCKETS 1` table this
    happens for every redirect that does not pick the single owning backend;
    with more backends than
    buckets the load fails whenever the entry backend owns none of the
    partition's tablets.
    
    Observed on a cloud regression run: 8 stream loads of single bucket
    random tables failed with the
    error above, while the same suite passed whenever the entry backend
    happened to be the tablet
    owner.
    
    ### What this PR changes
    
    * Resolve the adaptive routing for the backend that issued the request
    (`TStreamLoadPutRequest.backendId`, which the BE always sets for its own
    stream load) right before
    the plan is returned, reusing the existing compute/apply helpers. The
    sink then routes each
    partition to a backend it can actually reach, and the receiver side gets
    a matching local bucket
      sequence.
    * When the request carries no usable backend id (old clients, multi
    table stream load), no
    assignment consistent with the receiver side can be computed. In that
    case adaptive mode is turned
    off for the plan so the BE uses the legacy per-batch tablet routing,
    instead of letting it guess
    the bucket owner. This keeps the pre-#62661 behavior instead of failing.
    
    Routine load picks the executing BE after the plan is generated, so it
    is not covered here and
    still needs the sink backend id at planning time;
    `load_to_single_tablet=true` remains unaffected.
---
 .../org/apache/doris/load/StreamLoadHandler.java   |  74 ++++++++++
 .../load/routineload/kafka/KafkaTaskInfo.java      |   3 +
 .../load/routineload/kinesis/KinesisTaskInfo.java  |   3 +
 .../apache/doris/load/StreamLoadHandlerTest.java   | 156 +++++++++++++++++++++
 .../test_adaptive_random_bucket_routine_load.out   |  23 +++
 ...test_adaptive_random_bucket_routine_load.groovy | 139 ++++++++++++++++++
 .../test_adaptive_random_bucket_stream_load.groovy |  84 +++++++++++
 7 files changed, 482 insertions(+)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/StreamLoadHandler.java 
b/fe/fe-core/src/main/java/org/apache/doris/load/StreamLoadHandler.java
index f2002c011eb..e4fe09598fd 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/load/StreamLoadHandler.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/load/StreamLoadHandler.java
@@ -41,10 +41,14 @@ import 
org.apache.doris.nereids.load.NereidsCloudStreamLoadPlanner;
 import org.apache.doris.nereids.load.NereidsStreamLoadPlanner;
 import org.apache.doris.nereids.load.NereidsStreamLoadTask;
 import org.apache.doris.planner.GroupCommitPlanner;
+import org.apache.doris.planner.OlapTableSink;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.service.ExecuteEnv;
 import org.apache.doris.system.Backend;
 import org.apache.doris.system.SystemInfoService;
+import org.apache.doris.thrift.TDataSinkType;
+import org.apache.doris.thrift.TOlapTablePartitionParam;
+import org.apache.doris.thrift.TOlapTableSink;
 import org.apache.doris.thrift.TPipelineFragmentParams;
 import org.apache.doris.thrift.TStreamLoadPutRequest;
 import org.apache.doris.thrift.TStreamLoadPutResult;
@@ -60,6 +64,7 @@ import org.apache.logging.log4j.Logger;
 
 import java.security.SecureRandom;
 import java.util.List;
+import java.util.Map;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.stream.Collectors;
@@ -285,6 +290,7 @@ public class StreamLoadHandler {
                     ? multiTableFragmentInstanceIdIndex.getAndIncrement() : 0;
             TPipelineFragmentParams result = null;
             result = planner.plan(streamLoadTask.getId(), index);
+            assignAdaptiveRandomBucket(result);
             result.setTableName(table.getName());
             
result.query_options.setFeProcessUuid(ExecuteEnv.getInstance().getProcessUUID());
             result.setIsMowTable(table.getEnableUniqueKeyMergeOnWrite());
@@ -321,4 +327,72 @@ public class StreamLoadHandler {
     public List getFragmentParams() {
         return fragmentParams;
     }
+
+    /**
+     * Stream load sends the whole plan to the BE that issued this request, so 
the sink routing
+     * has to be resolved for exactly that BE. Without it the BE only knows 
the sink is in
+     * adaptive random bucket mode, assumes it owns every bucket and fails with
+     * "unknown partition channel" as soon as a partition has no tablet on it.
+     */
+    void assignAdaptiveRandomBucket(TPipelineFragmentParams params) {
+        TOlapTableSink sink = getOlapTableSink(params);
+        if (!OlapTableSink.shouldAssignAdaptiveRandomBucket(sink)) {
+            return;
+        }
+        long sinkBackendId = request.isSetBackendId() ? request.getBackendId() 
: -1L;
+        assignAdaptiveRandomBucket(params, sinkBackendId, request.getDb(), 
request.getTbl());
+    }
+
+    public static void assignAdaptiveRandomBucket(TPipelineFragmentParams 
params, long sinkBackendId,
+            String dbName, String tableName) {
+        TOlapTableSink sink = getOlapTableSink(params);
+        if (!OlapTableSink.shouldAssignAdaptiveRandomBucket(sink)) {
+            return;
+        }
+        if (sinkBackendId <= 0 || !sink.isSetLocation() || sink.getLocation() 
== null) {
+            // Old clients do not report the executing BE, and without the 
sink backend id or the
+            // tablet locations no assignment consistent with the receiver 
side can be computed
+            // here. Fall back to the non-adaptive per-batch routing, which 
never depends on the
+            // bucket owner. This is a normal compatibility fallback, not an 
error.
+            LOG.info("disable adaptive random bucket, stream load sink backend 
id is {}, db={}, table={}",
+                    sinkBackendId, dbName, tableName);
+            sink.unsetEnableAdaptiveRandomBucket();
+            return;
+        }
+        int sinkInstanceNum = Math.max(params.getLocalParamsSize(), 1);
+        Map<Long, Map<Long, OlapTableSink.AdaptiveBucketAssignment>> 
assignments =
+                OlapTableSink.computeAdaptiveRandomBucketAssignments(
+                        Lists.newArrayList(sinkBackendId), 
sink.getPartition().getPartitions(),
+                        sink.getLocation().getTablets(), sinkInstanceNum);
+        Map<Long, OlapTableSink.AdaptiveBucketAssignment> partitionAssignments 
= assignments.get(sinkBackendId);
+        if ((partitionAssignments == null || partitionAssignments.isEmpty())
+                && !isFakePartitionParam(sink.getPartition())) {
+            // Nothing could be assigned, so the sink would stay in adaptive 
mode without any
+            // routing and the BE would guess the bucket owner again.
+            LOG.warn("disable adaptive random bucket, no partition could be 
assigned for backend {}, "
+                            + "db={}, table={}", sinkBackendId, dbName, 
tableName);
+            sink.unsetEnableAdaptiveRandomBucket();
+            return;
+        }
+        OlapTableSink.applyAdaptiveRandomBucketAssignments(
+                sink.getPartition().getPartitions(), partitionAssignments);
+    }
+
+    /**
+     * Auto partition tables with no partition yet start with a fake partition 
list that is only
+     * used for tablet locations; their real partitions are created and 
assigned while the load
+     * runs, so an empty plan time assignment is expected and must keep 
adaptive mode enabled.
+     */
+    private static boolean isFakePartitionParam(TOlapTablePartitionParam 
partitionParam) {
+        return partitionParam.isSetPartitionsIsFake() && 
partitionParam.isPartitionsIsFake();
+    }
+
+    private static TOlapTableSink getOlapTableSink(TPipelineFragmentParams 
params) {
+        if (params == null || !params.isSetFragment() || params.getFragment() 
== null
+                || !params.getFragment().isSetOutputSink() || 
params.getFragment().getOutputSink() == null
+                || params.getFragment().getOutputSink().getType() != 
TDataSinkType.OLAP_TABLE_SINK) {
+            return null;
+        }
+        return params.getFragment().getOutputSink().getOlapTableSink();
+    }
 }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaTaskInfo.java
 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaTaskInfo.java
index 98aabe63926..322a38151c6 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaTaskInfo.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaTaskInfo.java
@@ -25,6 +25,7 @@ import org.apache.doris.common.Config;
 import org.apache.doris.common.UserException;
 import org.apache.doris.common.util.DebugPointUtil;
 import org.apache.doris.common.util.DebugUtil;
+import org.apache.doris.load.StreamLoadHandler;
 import org.apache.doris.load.routineload.RLTaskTxnCommitAttachment;
 import org.apache.doris.load.routineload.RoutineLoadJob;
 import org.apache.doris.load.routineload.RoutineLoadManager;
@@ -187,6 +188,8 @@ public class KafkaTaskInfo extends RoutineLoadTaskInfo {
                 (OlapTable) 
db.getTableOrMetaException(routineLoadJob.getTableId(),
                 Table.TableType.OLAP), taskInfo);
         TPipelineFragmentParams tExecPlanFragmentParams = 
routineLoadJob.plan(planner, loadId, txnId);
+        StreamLoadHandler.assignAdaptiveRandomBucket(tExecPlanFragmentParams, 
beId,
+                routineLoadJob.getDbFullName(), routineLoadJob.getTableName());
         TPlanFragment tPlanFragment = tExecPlanFragmentParams.getFragment();
         tPlanFragment.getOutputSink().getOlapTableSink().setTxnId(txnId);
 
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisTaskInfo.java
 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisTaskInfo.java
index 9c225a51250..a646b17cc5b 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisTaskInfo.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kinesis/KinesisTaskInfo.java
@@ -24,6 +24,7 @@ import org.apache.doris.catalog.Table;
 import org.apache.doris.common.Config;
 import org.apache.doris.common.UserException;
 import org.apache.doris.common.util.DebugUtil;
+import org.apache.doris.load.StreamLoadHandler;
 import org.apache.doris.load.routineload.RoutineLoadJob;
 import org.apache.doris.load.routineload.RoutineLoadManager;
 import org.apache.doris.load.routineload.RoutineLoadTaskInfo;
@@ -201,6 +202,8 @@ public class KinesisTaskInfo extends RoutineLoadTaskInfo {
                 (OlapTable) 
db.getTableOrMetaException(routineLoadJob.getTableId(),
                 Table.TableType.OLAP), taskInfo);
         TPipelineFragmentParams tExecPlanFragmentParams = 
routineLoadJob.plan(planner, loadId, txnId);
+        StreamLoadHandler.assignAdaptiveRandomBucket(tExecPlanFragmentParams, 
beId,
+                routineLoadJob.getDbFullName(), routineLoadJob.getTableName());
         TPlanFragment tPlanFragment = tExecPlanFragmentParams.getFragment();
         tPlanFragment.getOutputSink().getOlapTableSink().setTxnId(txnId);
 
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/load/StreamLoadHandlerTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/load/StreamLoadHandlerTest.java
index 1b925f06e09..b8231f963cd 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/load/StreamLoadHandlerTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/load/StreamLoadHandlerTest.java
@@ -29,9 +29,20 @@ import org.apache.doris.mysql.privilege.Auth;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.system.Backend;
 import org.apache.doris.system.SystemInfoService;
+import org.apache.doris.thrift.TDataSink;
+import org.apache.doris.thrift.TDataSinkType;
+import org.apache.doris.thrift.TOlapTableIndexTablets;
+import org.apache.doris.thrift.TOlapTableLocationParam;
+import org.apache.doris.thrift.TOlapTablePartition;
+import org.apache.doris.thrift.TOlapTablePartitionParam;
+import org.apache.doris.thrift.TOlapTableSink;
+import org.apache.doris.thrift.TPipelineFragmentParams;
+import org.apache.doris.thrift.TPlanFragment;
 import org.apache.doris.thrift.TStreamLoadPutRequest;
 import org.apache.doris.thrift.TStreamLoadPutResult;
+import org.apache.doris.thrift.TTabletLocation;
 
+import com.google.common.collect.Lists;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 import org.mockito.MockedStatic;
@@ -145,6 +156,151 @@ public class StreamLoadHandlerTest {
         }
     }
 
+    @Test
+    public void testAssignAdaptiveRandomBucketRoutesPartitionToTabletOwner() {
+        // The load runs on backend 20, but the only tablet of the partition 
lives on backend 10,
+        // so the sink has to route the partition to backend 10 instead of 
assuming self ownership.
+        TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+        request.setBackendId(20L);
+        StreamLoadHandler handler = new StreamLoadHandler(
+                request, null, new TStreamLoadPutResult(), "127.0.0.1");
+
+        TPipelineFragmentParams params = 
buildStreamLoadParams(createSingleBucketSink());
+        handler.assignAdaptiveRandomBucket(params);
+
+        TOlapTablePartition partition = getSinglePartition(params);
+        Assertions.assertEquals(10L, partition.getBucketBeId());
+        Assertions.assertEquals(0L, partition.getLoadTabletIdx());
+        Assertions.assertEquals(Lists.newArrayList(0), 
partition.getLocalBucketSeqs());
+        Assertions.assertEquals(10L, 
partition.getIndexes().get(0).getBucketBeId());
+        Assertions.assertEquals(Lists.newArrayList(0), 
partition.getIndexes().get(0).getLocalBucketSeqs());
+    }
+
+    @Test
+    public void 
testAssignAdaptiveRandomBucketKeepsLocalBucketWhenSinkOwnsTablet() {
+        TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+        request.setBackendId(10L);
+        StreamLoadHandler handler = new StreamLoadHandler(
+                request, null, new TStreamLoadPutResult(), "127.0.0.1");
+
+        TPipelineFragmentParams params = 
buildStreamLoadParams(createSingleBucketSink());
+        handler.assignAdaptiveRandomBucket(params);
+
+        TOlapTablePartition partition = getSinglePartition(params);
+        Assertions.assertEquals(10L, partition.getBucketBeId());
+        Assertions.assertEquals(Lists.newArrayList(0), 
partition.getLocalBucketSeqs());
+    }
+
+    @Test
+    public void testAssignAdaptiveRandomBucketFallsBackWithoutBackendId() {
+        // An old client does not report the executing backend, so no 
assignment consistent with the
+        // receiver side can be computed. Adaptive mode must be turned off to 
keep the legacy routing.
+        TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+        request.setDb("test_db");
+        request.setTbl("test_tbl");
+        StreamLoadHandler handler = new StreamLoadHandler(
+                request, null, new TStreamLoadPutResult(), "127.0.0.1");
+
+        TPipelineFragmentParams params = 
buildStreamLoadParams(createSingleBucketSink());
+        handler.assignAdaptiveRandomBucket(params);
+
+        TOlapTableSink sink = 
params.getFragment().getOutputSink().getOlapTableSink();
+        Assertions.assertFalse(sink.isSetEnableAdaptiveRandomBucket());
+        TOlapTablePartition partition = getSinglePartition(params);
+        Assertions.assertFalse(partition.isSetBucketBeId());
+        
Assertions.assertFalse(partition.getIndexes().get(0).isSetBucketBeId());
+    }
+
+    @Test
+    public void 
testAssignAdaptiveRandomBucketFallsBackWhenNothingCanBeAssigned() {
+        // A partition the FE cannot assign (here: without a load tablet 
index) would leave the sink
+        // in adaptive mode without any routing, so adaptive mode has to be 
turned off.
+        TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+        request.setBackendId(20L);
+        StreamLoadHandler handler = new StreamLoadHandler(
+                request, null, new TStreamLoadPutResult(), "127.0.0.1");
+
+        TOlapTableSink sink = createSingleBucketSink();
+        sink.getPartition().getPartitions().get(0).unsetLoadTabletIdx();
+        TPipelineFragmentParams params = buildStreamLoadParams(sink);
+        handler.assignAdaptiveRandomBucket(params);
+
+        Assertions.assertFalse(sink.isSetEnableAdaptiveRandomBucket());
+    }
+
+    @Test
+    public void testAssignAdaptiveRandomBucketKeepsAdaptiveForFakePartitions() 
{
+        // Auto partition tables get their real partitions, and their 
assignments, while the load
+        // runs, so an empty plan time assignment must keep adaptive mode 
enabled.
+        TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+        request.setBackendId(20L);
+        StreamLoadHandler handler = new StreamLoadHandler(
+                request, null, new TStreamLoadPutResult(), "127.0.0.1");
+
+        TOlapTableSink sink = createSingleBucketSink();
+        sink.getPartition().setPartitionsIsFake(true);
+        sink.getPartition().getPartitions().get(0).unsetLoadTabletIdx();
+        TPipelineFragmentParams params = buildStreamLoadParams(sink);
+        handler.assignAdaptiveRandomBucket(params);
+
+        Assertions.assertTrue(sink.isSetEnableAdaptiveRandomBucket());
+        Assertions.assertTrue(sink.isEnableAdaptiveRandomBucket());
+    }
+
+    @Test
+    public void testAssignAdaptiveRandomBucketIgnoresNonOlapTableSink() {
+        TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+        request.setBackendId(20L);
+        StreamLoadHandler handler = new StreamLoadHandler(
+                request, null, new TStreamLoadPutResult(), "127.0.0.1");
+
+        TPipelineFragmentParams params = new TPipelineFragmentParams();
+        TPlanFragment fragment = new TPlanFragment();
+        fragment.setOutputSink(new TDataSink(TDataSinkType.DATA_STREAM_SINK));
+        params.setFragment(fragment);
+
+        handler.assignAdaptiveRandomBucket(params);
+        
Assertions.assertNull(params.getFragment().getOutputSink().getOlapTableSink());
+    }
+
+    /**
+     * Creates a random distribution sink with one bucket, whose only tablet 
is located on backend 10.
+     */
+    private static TOlapTableSink createSingleBucketSink() {
+        TOlapTablePartition partition = new TOlapTablePartition();
+        partition.setId(1000L);
+        partition.setNumBuckets(1);
+        partition.setLoadTabletIdx(0);
+        partition.addToIndexes(new TOlapTableIndexTablets(1L, 
Lists.newArrayList(100L)));
+        TOlapTablePartitionParam partitionParam = new 
TOlapTablePartitionParam();
+        partitionParam.addToPartitions(partition);
+
+        TOlapTableLocationParam locationParam = new TOlapTableLocationParam();
+        locationParam.addToTablets(new TTabletLocation(100L, 
Lists.newArrayList(10L)));
+
+        TOlapTableSink sink = new TOlapTableSink();
+        sink.setEnableAdaptiveRandomBucket(true);
+        sink.setLoadToSingleTablet(false);
+        sink.setPartition(partitionParam);
+        sink.setLocation(locationParam);
+        return sink;
+    }
+
+    private static TPipelineFragmentParams 
buildStreamLoadParams(TOlapTableSink sink) {
+        TPipelineFragmentParams params = new TPipelineFragmentParams();
+        TPlanFragment fragment = new TPlanFragment();
+        TDataSink tDataSink = new TDataSink(TDataSinkType.OLAP_TABLE_SINK);
+        tDataSink.setOlapTableSink(sink);
+        fragment.setOutputSink(tDataSink);
+        params.setFragment(fragment);
+        return params;
+    }
+
+    private static TOlapTablePartition 
getSinglePartition(TPipelineFragmentParams params) {
+        return params.getFragment().getOutputSink().getOlapTableSink()
+                .getPartition().getPartitions().get(0);
+    }
+
     private Backend createBackend(long id, String host) {
         Backend backend = new Backend(id, host, 9050);
         backend.setAlive(true);
diff --git 
a/regression-test/data/load_p0/routine_load/test_adaptive_random_bucket_routine_load.out
 
b/regression-test/data/load_p0/routine_load/test_adaptive_random_bucket_routine_load.out
new file mode 100644
index 00000000000..3bd9915d837
--- /dev/null
+++ 
b/regression-test/data/load_p0/routine_load/test_adaptive_random_bucket_routine_load.out
@@ -0,0 +1,23 @@
+-- This file is automatically generated. You should know what you did if you 
want to edit this
+-- !first_batch --
+1      value_1
+2      value_2
+3      value_3
+4      value_4
+5      value_5
+6      value_6
+
+-- !second_batch --
+1      value_1
+10     value_10
+11     value_11
+12     value_12
+2      value_2
+3      value_3
+4      value_4
+5      value_5
+6      value_6
+7      value_7
+8      value_8
+9      value_9
+
diff --git 
a/regression-test/suites/load_p0/routine_load/test_adaptive_random_bucket_routine_load.groovy
 
b/regression-test/suites/load_p0/routine_load/test_adaptive_random_bucket_routine_load.groovy
new file mode 100644
index 00000000000..7099ce34797
--- /dev/null
+++ 
b/regression-test/suites/load_p0/routine_load/test_adaptive_random_bucket_routine_load.groovy
@@ -0,0 +1,139 @@
+// 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.
+
+import org.apache.doris.regression.suite.ClusterOptions
+import org.apache.doris.regression.util.RoutineLoadTestUtils
+import org.apache.kafka.clients.admin.AdminClient
+import org.apache.kafka.clients.admin.NewTopic
+import org.apache.kafka.clients.producer.ProducerConfig
+import org.apache.kafka.clients.producer.ProducerRecord
+
+suite("test_adaptive_random_bucket_routine_load", "docker") {
+    if (!RoutineLoadTestUtils.isKafkaTestEnabled(context)) {
+        return
+    }
+
+    def kafkaBroker = RoutineLoadTestUtils.getKafkaBroker(context)
+    def topic = 
"test_adaptive_random_bucket_routine_load_${UUID.randomUUID()}".toString()
+    def job = "test_adaptive_random_bucket_routine_load_job"
+    def options = new ClusterOptions()
+    options.setFeNum(1)
+    options.setBeNum(2)
+    options.cloudMode = true
+    options.feConfigs += [
+        'enable_adaptive_random_bucket_load=true',
+        'max_routine_load_task_num_per_be=1'
+    ]
+
+    docker(options) {
+        sql "DROP TABLE IF EXISTS test_adaptive_random_bucket_routine_load"
+        // Two concurrent Kafka tasks must use different BEs. With only one 
tablet,
+        // one task has to route to a BE other than the one executing the task.
+        sql """
+            CREATE TABLE test_adaptive_random_bucket_routine_load (
+                k INT NOT NULL,
+                v STRING
+            )
+            DUPLICATE KEY(k)
+            DISTRIBUTED BY RANDOM BUCKETS 1
+            PROPERTIES ("replication_allocation" = "tag.location.default: 1")
+        """
+
+        def adminProps = new Properties()
+        adminProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBroker)
+        def adminClient = AdminClient.create(adminProps)
+        try {
+            adminClient.createTopics([new NewTopic(topic, 2, (short) 
1)]).all().get()
+            try {
+                def producer = 
RoutineLoadTestUtils.createKafkaProducer(kafkaBroker)
+                try {
+                    def sendBatch = { int firstId ->
+                        for (int partition = 0; partition < 2; partition++) {
+                            for (int i = 0; i < 3; i++) {
+                                int id = firstId + partition * 3 + i
+                                producer.send(new ProducerRecord<>(topic, 
partition, null,
+                                        "${id},value_${id}".toString())).get()
+                            }
+                        }
+                        producer.flush()
+                    }
+
+                    sendBatch(1)
+                    sql """
+                        CREATE ROUTINE LOAD ${job} ON 
test_adaptive_random_bucket_routine_load
+                        COLUMNS TERMINATED BY ","
+                        PROPERTIES (
+                            "desired_concurrent_number" = "2",
+                            "max_batch_interval" = "5",
+                            "load_to_single_tablet" = "false"
+                        )
+                        FROM KAFKA (
+                            "kafka_broker_list" = "${kafkaBroker}",
+                            "kafka_topic" = "${topic}",
+                            "kafka_partitions" = "0,1",
+                            "kafka_offsets" = 
"OFFSET_BEGINNING,OFFSET_BEGINNING",
+                            "property.enable.partition.eof" = "false"
+                        )
+                    """
+                    try {
+                        def checkJob = {
+                            def state = sql_return_maparray "SHOW ROUTINE LOAD 
FOR ${job}"
+                            logger.info("Routine load status: ${state}")
+                            assertTrue(state[0].State in ["NEED_SCHEDULE", 
"RUNNING"],
+                                    "Routine load stopped unexpectedly: 
${state}")
+                            // Do not let a retry on the tablet owner hide a 
routing failure.
+                            
assertFalse(state[0].OtherMsg.toString().contains("unknown partition channel"),
+                                    "Adaptive routing failed: ${state}")
+                        }
+
+                        // Disable early EOF above so both tasks stay active 
for the batch interval,
+                        // allowing their BE assignments to be observed before 
they are renewed.
+                        awaitUntil(60) {
+                            checkJob()
+                            def tasks = sql_return_maparray "SHOW ROUTINE LOAD 
TASK WHERE JobName = '${job}'"
+                            def backendIds = tasks.collect { it.BeId as long 
}.findAll { it > 0 }.toSet()
+                            backendIds.size() == 2
+                        }
+
+                        def waitForRows = { int expected ->
+                            awaitUntil(120) {
+                                checkJob()
+                                def rows = sql "SELECT count(*) FROM 
test_adaptive_random_bucket_routine_load"
+                                rows[0][0] == expected
+                            }
+                        }
+                        waitForRows(6)
+                        order_qt_first_batch "SELECT k, v FROM 
test_adaptive_random_bucket_routine_load"
+
+                        // Renewed tasks must receive the adaptive assignments 
as well.
+                        sendBatch(7)
+                        waitForRows(12)
+                        order_qt_second_batch "SELECT k, v FROM 
test_adaptive_random_bucket_routine_load"
+                    } finally {
+                        sql "STOP ROUTINE LOAD FOR ${job}"
+                    }
+                } finally {
+                    producer.close()
+                }
+            } finally {
+                adminClient.deleteTopics([topic]).all().get()
+            }
+        } finally {
+            adminClient.close()
+        }
+    }
+}
diff --git 
a/regression-test/suites/load_p0/stream_load/test_adaptive_random_bucket_stream_load.groovy
 
b/regression-test/suites/load_p0/stream_load/test_adaptive_random_bucket_stream_load.groovy
new file mode 100644
index 00000000000..c07de6327fa
--- /dev/null
+++ 
b/regression-test/suites/load_p0/stream_load/test_adaptive_random_bucket_stream_load.groovy
@@ -0,0 +1,84 @@
+// 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.
+
+suite("test_adaptive_random_bucket_stream_load", "p0,nonConcurrent") {
+    if (!isCloudMode()) {
+        return
+    }
+
+    def enableAdaptiveRandomBucketConfig =
+            sql """ ADMIN SHOW FRONTEND CONFIG LIKE 
'enable_adaptive_random_bucket_load'; """
+    String oldEnableAdaptiveRandomBucket = 
enableAdaptiveRandomBucketConfig[0][1]
+
+    try {
+        sql """ ADMIN SET FRONTEND CONFIG 
('enable_adaptive_random_bucket_load' = 'true') """
+        sql """ DROP TABLE IF EXISTS test_adaptive_random_bucket_stream_load 
"""
+        // With a single bucket the partition has exactly one tablet, so all 
but one entry backend
+        // do not own it and have to route the partition to the owning backend.
+        sql """
+            CREATE TABLE test_adaptive_random_bucket_stream_load (
+                k int NOT NULL,
+                v string
+            )
+            DUPLICATE KEY(k)
+            DISTRIBUTED BY RANDOM BUCKETS 1
+            PROPERTIES (
+                "replication_allocation" = "tag.location.default: 1"
+            )
+        """
+
+        Map<String, String> backendIps = [:]
+        Map<String, String> backendHttpPorts = [:]
+        getBackendIpHttpPort(backendIps, backendHttpPorts)
+        assertTrue(backendIps.size() > 1, "adaptive random bucket routing 
needs at least two backends")
+
+        int rowsPerLoad = 10
+        int totalRows = 0
+        backendIps.each { backendId, backendIp ->
+            StringBuilder data = new StringBuilder()
+            for (int i = 0; i < rowsPerLoad; i++) {
+                data.append("${totalRows + i},value_${totalRows + i}\n")
+            }
+            totalRows += rowsPerLoad
+            // The entry backend is the sink backend. Only one of them owns 
the single tablet, the
+            // others must send the partition to the owning backend instead of 
failing with
+            // "unknown partition channel".
+            streamLoad {
+                table "test_adaptive_random_bucket_stream_load"
+                directToBe backendIp, backendHttpPorts.get(backendId) as int
+                set 'column_separator', ','
+                inputText data.toString()
+                time 60000
+
+                check { result, exception, startTime, endTime ->
+                    if (exception != null) {
+                        throw exception
+                    }
+                    def json = parseJson(result)
+                    
assertTrue(json.Status.toString().equalsIgnoreCase("success"), "load failed: 
${result}")
+                    assertEquals(rowsPerLoad, 
json.NumberLoadedRows.toString().toInteger(), "load failed: ${result}")
+                }
+            }
+            sql "sync"
+        }
+
+        def count = sql "SELECT count(*) FROM 
test_adaptive_random_bucket_stream_load"
+        assertEquals(totalRows, count[0][0] as int)
+    } finally {
+        sql """ ADMIN SET FRONTEND CONFIG 
('enable_adaptive_random_bucket_load' = '${oldEnableAdaptiveRandomBucket}') """
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to