OIiveirra commented on code in PR #68791:
URL: https://github.com/apache/doris/pull/68791#discussion_r4232924816


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/scan/FederationBackendPolicy.java:
##########
@@ -280,12 +297,15 @@ public Multimap<Backend, Split> 
computeScanRangeAssignment(List<Split> splits) t
                     }
                     case RANDOM: {
                         randomCandidates.reset();
-                        candidateNodes = 
selectNodes(Config.split_assigner_min_random_candidate_num, randomCandidates);
+                        candidateNodes = consistentHashSpreadNum == 0 ? 
backends
+                                : 
selectNodes(Config.split_assigner_min_random_candidate_num, randomCandidates);
                         break;
                     }
                     case CONSISTENT_HASHING: {
                         candidateNodes = consistentHash.getNode(split,
-                                
Config.split_assigner_min_consistent_hash_candidate_num);
+                                consistentHashSpreadNum == 0 ? backends.size()

Review Comment:
   已修复。spread 候选数达到当前 eligible BE 总数时,现在直接复用已过滤的 `backends` 列表;仅在候选数小于 BE 数量时遍历 
consistent-hash ring。`spread_num = 0` 仍直接使用全部候选 BE。修复已包含在提交 `886cc62894c`。



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/scan/FederationBackendPolicy.java:
##########
@@ -280,12 +297,15 @@ public Multimap<Backend, Split> 
computeScanRangeAssignment(List<Split> splits) t
                     }
                     case RANDOM: {
                         randomCandidates.reset();
-                        candidateNodes = 
selectNodes(Config.split_assigner_min_random_candidate_num, randomCandidates);
+                        candidateNodes = consistentHashSpreadNum == 0 ? 
backends
+                                : 
selectNodes(Config.split_assigner_min_random_candidate_num, randomCandidates);
                         break;
                     }
                     case CONSISTENT_HASHING: {
                         candidateNodes = consistentHash.getNode(split,
-                                
Config.split_assigner_min_consistent_hash_candidate_num);
+                                consistentHashSpreadNum == 0 ? backends.size()
+                                        : isSpreadEnabled() ? 
Math.min(consistentHashSpreadNum, backends.size())

Review Comment:
   已修复。为 Paimon、Trino、Fluss/FlussLake、ADBC 和 Remote Doris split 增加稳定的 connector 
split identity,并由 `PluginDrivenSplit` 传递到调度 Hash。pathless ranges 不再共享虚拟路径作为唯一 
Hash 输入。修复已包含在提交 `fe83abb6008`。



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/scan/FederationBackendPolicy.java:
##########
@@ -301,17 +321,23 @@ public Multimap<Backend, Split> 
computeScanRangeAssignment(List<Split> splits) t
                 throw new 
UserException(SystemInfoService.NO_SCAN_NODE_BACKEND_AVAILABLE_MSG);
             }
 
-            Backend selectedBackend = chooseNodeForSplit(candidateNodes);
-            List<Backend> alternativeBackends = new 
ArrayList<>(candidateNodes);
-            alternativeBackends.remove(selectedBackend);
-            split.setAlternativeHosts(
-                    alternativeBackends.stream().map(each -> 
each.getHost()).collect(Collectors.toList()));
+            Backend selectedBackend = isSpreadEnabled() && 
split.isRemotelyAccessible()
+                    ? chooseNodeForSpread(candidateNodes) : 
chooseNodeForSplit(candidateNodes);
+            // Alternative hosts are used only by global redistribution, which 
spread mode disables.
+            if (!isSpreadEnabled()) {
+                List<Backend> alternativeBackends = new 
ArrayList<>(candidateNodes);
+                alternativeBackends.remove(selectedBackend);
+                split.setAlternativeHosts(
+                        alternativeBackends.stream().map(each -> 
each.getHost()).collect(Collectors.toList()));
+            }
             assignment.put(selectedBackend, split);
             assignedWeightPerBackend.put(selectedBackend,
                     assignedWeightPerBackend.get(selectedBackend) + 
split.getSplitWeight().getRawValue());
         }
 
-        if (enableSplitsRedistribution) {
+        // Global redistribution can move a split outside its hash candidates 
or its locality constraints.
+        // The spread mode balances weights within each split's candidates 
during initial assignment instead.
+        if (enableSplitsRedistribution && !isSpreadEnabled()) {

Review Comment:
   已修复。spread 模式下,Host hint 只作为本地偏好;当偏好 BE 的已分配权重高于全局最小权重时,该 split 会进入 spread 
调度流程,从候选 BE 中重新选择最小权重节点。非 remote 的强制 Host 约束保持不变。修复已包含在 `fe83abb6008`。



##########
fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java:
##########
@@ -1486,6 +1486,17 @@ public enum IgnoreSplitType {
             description = "Use consistent hashing to split the appearance for 
external scan")
     public boolean useConsistentHashForExternalScan = false;
 
+    public static final String EXTERNAL_SCAN_CONSISTENT_HASH_SPREAD_NUM = 
"external_scan_consistent_hash_spread_num";
+    @VarAttrDef.VarAttr(name = EXTERNAL_SCAN_CONSISTENT_HASH_SPREAD_NUM,
+            checker = "checkExternalScanConsistentHashSpreadNum", needForward 
= true,

Review Comment:
   已修复。`use_consistent_hash_for_external_scan` 已设置 `needForward = true`,并在 
`VariableMgrTest` 中增加 forwarding 验证,确保 follower 转发 session 时 selector 与 spread 
count 一起生效。修复已包含在 `813826a7feb` 及后续分支提交中。



##########
regression-test/suites/external_table_p0/tvf/test_external_scan_consistent_hash_spread.groovy:
##########
@@ -0,0 +1,88 @@
+// 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_external_scan_consistent_hash_spread", "p0,external") {
+    String ak = getS3AK()
+    String sk = getS3SK()
+    String endpoint = getS3Endpoint()
+    String region = getS3Region()
+    String bucket = getS3BucketName()
+    String pathStyle = getS3Provider().equalsIgnoreCase("S3") ? "true" : 
"false"
+    String prefix = 
"s3://${bucket}/test_external_scan_consistent_hash_spread/${UUID.randomUUID()}/part_"
+
+    sql "drop table if exists test_external_scan_consistent_hash_spread"
+    sql """
+        create table test_external_scan_consistent_hash_spread (id int, value 
string)
+        distributed by hash(id) buckets 1
+        properties ("replication_num" = "1")
+    """
+    sql "insert into test_external_scan_consistent_hash_spread values (1, 
'alpha'), (2, null), (3, 'gamma')"
+    sql """
+        select id, value from test_external_scan_consistent_hash_spread order 
by id
+        into outfile "${prefix}single_" format as parquet
+        properties (
+            "s3.endpoint" = "${endpoint}", "s3.region" = "${region}",
+            "s3.access_key" = "${ak}", "s3.secret_key" = "${sk}",
+            "s3.path_style_access" = "${pathStyle}"

Review Comment:
   已修复。两个 OUTFILE 导出块均改用 S3 filesystem 支持的 `use_path_style` 属性,避免 
path-style-only endpoint 在回归测试的扫描验证前失败。修复已包含在提交 `886cc62894c`。



-- 
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