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]

Reply via email to