github-actions[bot] commented on code in PR #68006:
URL: https://github.com/apache/doris/pull/68006#discussion_r4120289721


##########
fe/fe-core/src/main/java/org/apache/doris/load/StreamLoadHandler.java:
##########
@@ -321,4 +327,72 @@ public void generatePlan() throws UserException, 
MetaNotFoundException {
     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(

Review Comment:
   [P2] Avoid per-partition INFO logging on every adaptive load plan. With no 
explicit partition filter, `OlapTableSink.init` includes every table partition. 
This new call reaches `computeAdaptiveRandomBucketAssignments` and 
`applyAdaptiveRandomBucketAssignments`, which each log an INFO record per 
assigned partition. A 1,000-partition random table therefore produces roughly 
2,000 detailed FE log records for one successful stream load, and the same cost 
recurs for each Kafka/Kinesis task. This is distinct from the single 
missing-backend fallback log discussed in the existing thread. Please make the 
detailed assignment messages DEBUG-level or emit one bounded summary per plan 
so frequent loads do not flood FE logs.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to