OIiveirra opened a new pull request, #68791:
URL: https://github.com/apache/doris/pull/68791
### What problem does this PR solve?
Issue Number: N/A
Related PR: #68726
This PR is submitted independently against master. It does not include the
related PR's commits or change that PR; merging it is not a prerequisite for
reviewing this scheduling change. Joint validation results below are
explicitly
identified as a local integration snapshot.
Problem Summary:
Independent queries for the same remote external file split can repeatedly
choose the same backend under consistent hashing. The hash candidate set is
stable, each policy starts with zero assigned weight, and the existing tie
rule
chooses the same candidate. Increasing the existing candidate count alone
does
not remove this hotspot.
This change adds an opt-in `external_scan_consistent_hash_spread_num` session
setting. Its default value of `1` preserves the original candidate count, tie
rule and redistribution. Values greater than `1` apply only when the scan
already uses consistent hashing:
```text
remote split without preferred local backend
-> first N distinct eligible hash candidates, capped by eligible BE count
-> least policy-assigned weight
-> uniformly random choice among tied candidates
-> assign the split once
```
The new mode retains query-local weight history across batches and honors
local preference, mandatory Host constraints and compute-group eligibility.
Global redistribution is disabled only in the new mode because it could move
a
split outside its hash candidates or mandatory locality. With many splits,
balancing is consequently limited to each split's candidates. No global query
counter, CPU monitoring, COUNT pushdown, DOP change or BE cache protocol is
added.
A local integration snapshot includes the unchanged related PR on both sides.
With the same FE binary, three unchanged BE binaries, and the same
million-row
single-file CSV workload at concurrency 64, changing only the spread setting
from `1` to `3` produced these three-round paired cache-off results:
| Query | Original QPS | Spread QPS | Pooled ratio | Per-round ratio range |
|---|---:|---:|---:|---:|
| COUNT | 45.056 | 132.789 | 2.9472x | 2.9241x–2.9874x |
| SUM | 13.278 | 38.772 | 2.9201x | 2.8887x–2.9574x |
Each measurement used a 10-second warmup and a 60-second window. QPS counts
correct queries completed in that window; all 12 measurements had zero query
errors. Complete-call latency also includes successful queries initiated in
the
window that completed after its end. The original hot BE consumed about 98%
of
its four-core affinity capacity. With spreading, all three BEs consumed about
90.6%–97.8% of their respective capacities.
These are measurements on one shared physical host with separate logical
Hosts
and disjoint, nonexclusive CPU affinities. The integrated FE is a local
snapshot,
not proof that the related PR has merged upstream. COUNT pushdown and the
other
related PR changes are identical on both sides; this PR adds only scheduling.
The measurements show approximately threefold throughput in this workload,
not a strict >=3.000x result, a universal guarantee, three-host acceptance or
customer validation. Earlier measurements based on master `c9c85e75539`,
without the related
PR, remain separately recorded: concurrency-32 pooled ratios were 2.4234x
(COUNT) and 2.4143x (SUM),
and concurrency-64 ratios were 2.5684x and 2.5843x.
Optional hardware counter collection failed to attach in eight groups; of
four
groups with numeric counters, two also reported attachment warnings. No
paired
IPC explanation is inferred. The QPS, CPU sampling, query results and runtime
identity checks are complete. Global scan counters have additional producers
and are not used as query-specific proof of execution counts.
Cache cost measurements showed one warmed backend copy in the original mode
and three in spread mode, with three times the initial remote read bytes. A
separate three-round warm-cache COUNT experiment on the same candidate FE
recorded a pooled 2.6973x ratio at concurrency 32, zero errors and zero
additional
remote reads. This cache-on result is separate from the cache-off scheduling
comparison.
### Release note
Add an optional `external_scan_consistent_hash_spread_num` session setting to
spread repeated external scan queries across consistent hash candidates. The
default preserves existing scheduling. Enabling spreading with file caching
can increase warmup reads and cache space on additional backends.
### Check List (For Author)
- Test
- [x] Regression test
- [x] Unit Test
- [x] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason
Validation performed:
- Standard FE build and Checkstyle passed.
- `FederationBackendPolicySpreadTest`: 19 tests passed.
- `FederationBackendPolicyTest`: 12 existing tests passed in the joint
snapshot.
- `VariableMgrTest#testExternalScanConsistentHashSpreadVariable` passed.
- `test_external_scan_consistent_hash_spread_variable` and
`test_external_scan_consistent_hash_spread` regression suites were
generated
and verified through the standard regression entry point using an
isolated
personal cluster and a self-generated three-row Parquet fixture in MinIO.
- In the joint snapshot, COUNT and SUM were each executed 100 times with
spread setting `1` and 100 times with `3` (400 correct queries). Every
query
had one split and one positive execution target; `1` kept the hotspot and
`3` selected all three BEs. Both sides read one million CSV rows; COUNT
used the same existing pushdown with zero deserialized cells, while SUM
deserialized one million cells.
- Manual capacity experiment: use a single million-row CSV file, execute
`COUNT(*)` and `SUM` over that same file with SQL/query/file caches
disabled
and consistent hashing enabled, compare original scheduling with explicit
spread count `3` at concurrency `32`, alternate A/B ordering for three
paired 60-second rounds, and record successful completions, errors,
complete latency, backend CPU, hashes and runtime identities.
- Cold-cache cost sampling and separate warm-cache COUNT pairs passed.
- A separate same-binary, cache-off concurrency-64 comparison completed
three paired rounds: pooled ratios 2.5684x (COUNT) and 2.5843x (SUM),
with zero query errors in all 12 measurements. A diagnostic runner
interruption and recovery is retained in its evidence. Global scan
counters include internal/other work and are not query-specific proof.
- Joint snapshot: standard FE build/Checkstyle passed; 39 focused tests
passed (12 existing backend policy, 4 COUNT pushdown, 3 remote scan job,
19 spread policy and 1 session variable).
- Both targeted regression suites were verified again on the joint snapshot
through the personal standard wrapper, each with 1 successful suite and
zero failed/fatal/skipped scripts. Health passed before and after, and
FE/BE PIDs and hashes were unchanged.
- Joint snapshot manual capacity test: same FE, spread `1` versus `3`,
concurrency `64`, same COUNT/SUM SQL and caches/settings on both sides,
three alternating pairs of 60-second windows per SQL. All 12 groups and
frozen-input/driver/settings/identity/timing checks passed. Results
above.
- This targeted local coverage does not claim full TeamCity or three-host
acceptance. No BE source changed, so BE UT/static analysis were not
rerun.
- Behavior changed:
- [ ] No.
- [x] Yes. When the new setting is greater than 1 and the scan uses
consistent hashing, equal-weight hash candidates are selected randomly
and global redistribution is disabled. Default scheduling is preserved.
- Does this need documentation?
- [ ] No.
- [x] Yes. Included in `docs/external-scan-consistent-hash-spread.md`.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
--
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]