github-actions[bot] commented on code in PR #68791:
URL: https://github.com/apache/doris/pull/68791#discussion_r4229303072
##########
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:
[P2] Forward the hash selector together with the new spread count.
`use_consistent_hash_for_external_scan` above has no `needForward`, so a SELECT
forwarded from a follower with file cache off, that flag on, and spread 2 sends
only the new integer. The master sees its default false selector and
`createBackendPolicy` chooses round-robin with spread 1, silently dropping
bounded hashing. Mark the selector for forwarding and test policy construction
from a forwarded session-variable map.
##########
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:
[P2] Avoid walking the full hash ring when all BEs are spread candidates.
`spreadNum == 0`, or a spread value >1 at least the eligible BE count, asks
`getNode` for every BE on every remote split. That method scans the
virtual-node TreeMap (256 replicas per BE by default) and allocates a
distinct-node set and list, even though `chooseNodeForSpread` only needs the
already-filtered `backends` list to compare weights and draw a uniform tie.
Large multi-file scans repeat this work for each split, including within the
lazy assignment lock. Use `backends` directly for this all-candidate case.
##########
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:
[P2] Use the recognized path-style key for the OUTFILE fixture. Both exports
set `s3.path_style_access`, but the S3 filesystem property binder accepts
`use_path_style` or `s3.path-style-access`; the misspelled key leaves the
default `use_path_style=false` in the BE sink properties. On a path-style-only
S3-compatible endpoint, the export fails before any spread query runs. Use
`use_path_style` here and in the second OUTFILE block (line 81), matching the
TVF read properties.
##########
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:
[P2] Give pathless ranges a distinct stable hash identity before bounding
their candidates. Paimon JNI creates separate ranges with different serialized
`paimon.split` values but no path, start or length; `PluginDrivenSplit` maps
each to `(connector://virtual, 0, -1)`, which is all `SplitHash` reads. With
100 target-size-or-larger ranges, 10 eligible BEs and spread 2, every range
gets the same two candidates and this mode skips redistribution, leaving eight
BEs idle; the old mode rebalanced this load. Trino and remote Doris have
similar placeholder keys. Use connector work-unit identity or retain a
balancing path for these ranges.
##########
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:
[P2] Keep spreading host-hinted remote splits when the preferred BE is
overloaded. The default local-preference pass assigns every remotely accessible
split with a matching single host before this branch. For 30 such splits on BE
A and three eligible BEs, `remainingSplits` is empty; skipping
`equateDistribution` leaves A with all 30 and B/C idle, whereas the original
path rebalances them. Trino and ES ranges can provide these hosts. Let
overloaded remote splits enter a constrained spread/rebalance decision while
preserving mandatory hosts for nonremote splits.
--
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]