NoahKusaba opened a new issue, #2520:
URL: https://github.com/apache/datafusion-ballista/issues/2520

   **Is your feature request related to a problem or challenge? Please describe 
what you are trying to do.**
   
   This tracks a performance review of the shuffle data path, executor runtime, 
scheduler control plane and client result path, looking only for throughput, 
latency and memory wins. The static findings were checked against the TPC-H 
SF1000 wall-clock runs in #2499 and the per-stage plans and task metrics in 
#2419. The full write-up, with plan excerpts and per-stage tables, is here: 
https://claude.ai/artifact/M1RdVxs9H9FgnYVsncomXF
   
   Line numbers refer to 187e21fc. **Impact** is the expected effect on query 
time or peak memory for shuffle-heavy workloads on a multi-executor cluster. 
**Effort:** S = hours, M = days, L = weeks. Rankings for revisions 1–2 are 
estimates from reading the code; no profiles were run. A11, A12, C15, C16 and 
C17 come from #2419's plan shapes and task-time distributions, and their code 
locations still need pinning down.
   
   **Describe the solution you'd like**
   
   Each item can be split into its own issue or PR when someone picks it up. 
Items marked with a PR already have one.
   
   | # | Finding | Impact | Effort | Status |
   |---|---|---|---|---|
   | A1 | Remote shuffle blocks fully buffered before first batch | High | M |  
|
   | A2 | Coarse tasks hold every vcore until the slowest partition finishes | 
High | M |  |
   | A3 | Full stage plan re-encoded and shipped per task | High | M |  |
   | A4 | Scheduler serialised through one event loop, uncoalesced revives | 
High | M |  |
   | A11 | Small filtered build sides don't stop the lineitem shuffle | High | 
M |  |
   | A5 | Executor memory sliced statically per vcore | Medium | M |  |
   | A6 | Sort-shuffle layout forces a barrier and copies spilled data | Medium 
| L |  |
   | A12 | Byte-range splits of large row groups leave one task doing the work 
| Medium | M |  |
   | A7 | Broadcast build sides re-fetched and rebuilt per task | Medium | M |  
|
   | A8 | Results fetched after job end, serially, through Flight | Medium | M 
|  |
   | A9 | Blocking I/O fans out to one OS thread per output partition | Medium 
| M |  |
   | A10 | Task placement ignores shuffle data locality | Low–Med | M |  |
   | C1 | Sort-shuffle spill and final write block tokio workers | High | S |  |
   | C15 | Probe side read as broadcast runs the join in one task | Medium | S 
| #2501 |
   | C16 | Date-partitioned scans list every partition directory | Medium | S | 
 |
   | C2 | Flight `do_get` re-compresses with hard-coded LZ4 | Medium | S |  |
   | C3 | Client result fetch: new connection per partition, 64 KiB windows | 
Medium | S |  |
   | C4 | Shuffle index and schema header re-read per fetch | Medium | S |  |
   | C5 | Passthrough shuffle writer writes small batches as-is | Medium | S | 
#2518 |
   | C6 | Sort-shuffle end-of-input encode doubles peak memory | Medium | M |  |
   | C7 | Each revive clones the job map and write-locks every graph | Medium | 
S |  |
   | C8 | Sort-shuffle local reads use an unbuffered `File` | Medium | S |  |
   | C9 | Local shuffle files read one at a time inside `poll_next` | Low–Med | 
S |  |
   | C10 | Remote fetch spawns one tokio task per location up front | Low–Med | 
S |  |
   | C11 | Every `PartitionLocation` carries full executor metadata | Low–Med | 
S |  |
   | C17 | Stage row metrics sum across operators and switch units | Low | S |  
|
   | C12 | `ballista_config()` clones and re-parses settings per call | Low | S 
|  |
   | C13 | View columns garbage-collected unconditionally | Low | S |  |
   | C14 | Pull-mode client polls job status every 50 ms | Low | S |  |
   | R1 | Scheduler re-read every Parquet footer per job (#2497) | — | — | 
Fixed in #2498 |
   
   ### SF1000 context
   
   - #2499: Ballista totals 652.5 s, 0.72× of Spark 3.5 (909.6 s) and 1.45× of 
Spark + Comet (450.0 s). 100.2 s of that was planning (R1, fixed in #2498); 
without it Ballista is about 1.21× of Comet.
   - Largest remaining gaps to Comet: q06 (11.5×, mostly R1 then C16), q14 
(4.2×, R1 + C16), q07 (3.8×, a 14.6 s max vs 1.4 s median task in stage 1), q20 
(3.4×, two single-task stages, A12).
   - Outlier stages in #2419: q08 s1 (full lineitem sort-shuffle to join 1.3M 
parts, task time 14.2 / 62.5 / 131.1 s min/median/max on a 1.01× input spread: 
A11, A2, A6, C1), q08 s5 (80M probe rows in one task: C15), q02 s6, q10 s0, q16 
s4 and q20 s5 (one task doing the whole scan: A12), q06 s0 (one-year filter 
over all 256 tasks: C16).
   
   ### Architectural
   
   #### A1. Remote shuffle blocks are fully buffered before the first batch is 
emitted
   **Impact:** High · **Effort:** M
   
   Every remote partition fetch does `stream.try_collect::<Vec<_>>()` before 
handing anything downstream, so network transfer and compute never overlap 
within a block. The governor charges compressed wire bytes while decoded 
batches stay resident, so peak memory is roughly the compression ratio × 
`max_bytes_in_flight` (48 MiB default), and none of it is registered with the 
task's `MemoryPool`. The buffering exists so a mid-body failure can be retried 
without emitting duplicates.
   
   **Suggested change:** stream batches as they decode and make retries 
resumable by skipping the batches already emitted per location (IPC batch order 
within a block is deterministic). If buffering must stay, buffer compressed 
bytes (or spill them) and charge a `MemoryReservation`.
   
   **Where:** 
[`shuffle_reader.rs:994-1022`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L994-L1022),
 
[`shuffle_reader.rs:896-924`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L896-L924)
   
   #### A2. Coarse tasks hold every vcore until the slowest partition finishes
   **Impact:** High · **Effort:** M
   
   `bind_one` packs up to `budget.vcores` partitions into one task, and 
`max_partitions_per_task` defaults to 0 (unbounded), so under the default 
`Bias` policy a stage is usually one task per executor. Vcores are only 
refunded when the whole task finishes, so one skewed partition keeps V−1 cores 
reserved and idle. A failure retries the whole slice, and there is no 
speculative execution. SF1000: q08 s1 ran 34 tasks with a 62.5 s median and 
131.1 s max on near-identical inputs.
   
   **Suggested change:** report progress per partition and refund each vcore as 
its drain finishes, or default `max_partitions_per_task` to 1–2 once A3 makes 
small tasks cheap. Add straggler detection and speculative re-launch; 
`(task_id, file_id)` naming already keeps duplicate attempts from colliding.
   
   **Where:** 
[`mod.rs:427-505`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L427-L505),
 
[`query_stage_scheduler.rs:483-498`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/scheduler_server/query_stage_scheduler.rs#L483-L498),
 
[`config.rs:472-480`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/config.rs#L472-L480)
   
   #### A3. Full stage plan is re-encoded and shipped for every task
   **Impact:** High · **Effort:** M
   
   `prepare_multi_task_definition` runs `restrict_plan_to_partitions` and a 
full `PhysicalPlanNode::try_from_physical_plan` + encode per task, and clones 
the session key/value props into every `MultiTaskDefinition`. Cost grows as 
tasks × plan size, on the same runtime as the event loop. At SF1000 every scan 
stage carries a `DataSourceExec` with 256 file groups listing every Parquet 
path and byte range, re-encoded for each task.
   
   **Suggested change:** encode the stage plan once per `(job, stage, 
attempt)`, cache it on executors, and send each task only its partition slice, 
reader locations and file groups. Send session props once per session.
   
   **Where:** 
[`task_manager.rs:877-945`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/state/task_manager.rs#L877-L945),
 
[`mod.rs:319-360`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/state/mod.rs#L319-L360)
   
   #### A4. Scheduler work is serialised through one event loop with 
uncoalesced revives
   **Impact:** High · **Effort:** M
   
   Every `TaskUpdating` batch posts a `ReviveOffers`, and each runs in full: it 
snapshots all running jobs and holds the global `available_vcores` mutex while 
awaiting each job graph's write lock in turn. AQE re-runs the whole 
physical-optimizer pipeline inside the loop when a stage resolves. Scheduling 
delay grows with jobs × stages × task completions, and jobs are visited in 
`HashMap` order.
   
   **Suggested change:** coalesce revives with a dirty flag, track pending-task 
counts atomically so idle jobs are skipped without locking, iterate jobs FIFO 
or fair-share, and move AQE replanning into a per-job task that posts back to 
the loop.
   
   **Where:** 
[`event_loop.rs:78-92`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/event_loop.rs#L78-L92),
 
[`query_stage_scheduler.rs:505-508`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/scheduler_server/query_stage_scheduler.rs#L505-L508),
 
[`query_stage_scheduler.rs:522-524`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/scheduler_server/query_stage_scheduler.rs#L522-L524),
 
[`memory.rs:73`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/memory.rs#L73),
 
[`mod.rs:538-545`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L538-L545),
 
[`planner.rs:356-384`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/state/aqe/planner.rs#L356-L384)
   
   #### A11. Small filtered build sides don't stop the full lineitem shuffle
   **Impact:** High · **Effort:** M
   
   In q08, q09 and q17 both sides of the part ⋈ lineitem join get an 
`ExchangeExec`, and `DynamicJoinSelectionExec` (`repartitioned=true`) only 
picks the strategy after both have run. The stages are independent, so by the 
time the part side reports 1.3M rows (q08, well under the 128 MiB broadcast 
threshold), all 6.0B lineitem rows are already being hash-shuffled: 131 s of 
q08's 164 s in #2419. q17 does it twice. No build-side filter reaches the 
lineitem scan.
   
   **Suggested change:** when one side has a selective filter over a small 
table, resolve it first and hold the other side's stage until its size is 
known; if it fits the broadcast threshold, replan to a broadcast join with no 
exchange on the lineitem side. Separately, push the resolved build side's keys 
(min/max, bloom or in-list) into the probe-side `DataSourceExec`.
   
   **Where:** #2419 `q08.md`, `q09.md`, `q17.md`, 
[`state/aqe/`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/state/aqe)
   
   #### A5. Executor memory is sliced statically per vcore
   **Impact:** Medium · **Effort:** M
   
   Each task gets its own `FairSpillPool` of `per_vcore × vcores_consumed` and 
spills when it fills its slice, even if the rest of the executor is idle. 
Collapse stages (final sort, SPM, global aggregate) run as one task with 
`vcores_consumed = 1`, so they get 1/V of executor memory. In #2419, 
sort-shuffle writers are capped at 256 MiB while buffering ~180M lineitem rows 
into 256 partitions.
   
   **Suggested change:** use one executor-wide pool with a guaranteed minimum 
per task and borrowable headroom, or at least size collapse tasks by the memory 
they need.
   
   **Where:** 
[`executor_process.rs:224-245`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/executor/src/executor_process.rs#L224-L245),
 
[`mod.rs:462-468`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L462-L468)
   
   #### A6. Sort-shuffle file layout forces a barrier and a second copy of 
spilled data
   **Impact:** Medium · **Effort:** L
   
   The coordinator writes nothing until all M input drains finish, so the final 
write is serial and on the critical path. Spilled bytes are written twice 
(spill file, then `std::io::copy` into `data.arrow`), and each (input × 
partition × spill) contribution is a separate IPC stream repeating the schema, 
which `MultiStreamPartitionStream` re-parses per segment. In #2419, q08 s1's 
`SortShuffleWriterExec` shows a 9.4× spread in task time on inputs within 1.36×.
   
   **Suggested change:** let the index point at a list of segments per 
partition (spill files plus the final flush) so it references data instead of 
copying it, or emit one IPC stream per partition without per-segment schemas. 
Either removes the M-input barrier.
   
   **Where:** 
[`writer.rs:1117-1186`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/writer.rs#L1117-L1186),
 
[`writer.rs:823-905`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/writer.rs#L823-L905),
 
[`multi_stream_reader.rs:95-108`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/multi_stream_reader.rs#L95-L108)
   
   #### A12. Byte-range splits of large row groups leave one task doing all the 
work
   **Impact:** Medium · **Effort:** M
   
   The SF1000 data uses ~512 MiB row groups. Scans split each file into byte 
ranges over 256 file groups, and a row group is read by the one range 
containing its midpoint, so with fewer row groups than ranges most ranges read 
nothing. Supplier stages put all 10M rows in one task (q02 s6, q15 s0, q16 s4, 
q20 s5); q10 s0 has a 0.2 s median and 17.1 s max.
   
   **Suggested change:** split scans on row-group boundaries using the footer 
metadata the scheduler already reads, capping the group count at the row-group 
count. For tables with one or two row groups, use page-level parallelism or 
schedule the single task first. Consider defaulting `coalesce.enabled` on at 
large scale factors.
   
   **Where:** #2419 `q02.md` s6, `q10.md` s0, `q15.md` s0, `q16.md` s4, 
`q20.md` s5
   
   #### A7. Broadcast build sides are re-fetched and rebuilt per task
   **Impact:** Medium · **Effort:** M
   
   A broadcast `ShuffleReaderExec` reads partition 0 with all locations. Across 
tasks on one executor there's no shared cache, so each task fetches the whole 
build side and builds its own hash table. Smaller tasks (A2) make this worse.
   
   **Suggested change:** add an executor-level broadcast cache keyed by `(job, 
stage, attempt)` holding the fetched batches (or built `JoinHashMap`) behind a 
once-cell, evicted on job cleanup.
   
   **Where:** 
[`shuffle_reader.rs:168-190`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L168-L190),
 
[`shuffle_reader.rs:374`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L374)
   
   #### A8. Results are fetched only after the job ends, serially, through 
Flight re-encoding
   **Impact:** Medium · **Effort:** M
   
   The client waits for `Successful`, then reads partitions one after another 
via `stream::iter(partitions).flatten()`. Each fetch creates a new 
`BallistaClient` with 64 KiB HTTP/2 windows and uses the Flight path, where the 
executor decodes and re-encodes with LZ4 (C2).
   
   **Suggested change:** stream final-stage partitions as they complete, with 
bounded concurrency over pooled connections using the configured windows, and 
use `IO_BLOCK_TRANSPORT` so executors send raw file bytes.
   
   **Where:** 
[`distributed_query.rs:586-640`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/distributed_query.rs#L586-L640),
 
[`distributed_query.rs:856-905`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/distributed_query.rs#L856-L905),
 
[`flight_service.rs:143-206`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/executor/src/flight_service.rs#L143-L206)
   
   #### A9. Blocking I/O fans out to one OS thread per output partition
   **Impact:** Medium · **Effort:** M
   
   `write_stream_to_disk` holds a `spawn_blocking` thread for the life of each 
output stream, and the passthrough writer runs K of them per task. IPC encoding 
and compression run on those threads, outside the `DedicatedExecutor`'s 
vcore-sized pool, so K = 200 means 200 OS threads doing CPU work, which breaks 
vcore accounting. Elsewhere, synchronous file I/O runs on async workers (C1, 
C9).
   
   **Suggested change:** encode and compress on the task runtime, then hand 
finished buffers to a bounded executor-wide I/O pool sized to the disks, and 
use that one policy for every file operation in the shuffle path.
   
   **Where:** 
[`utils.rs:200-270`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/utils.rs#L200-L270),
 
[`shuffle_writer.rs:586-625`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_writer.rs#L586-L625),
 
[`cpu_bound_executor.rs:100-115`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/executor/src/cpu_bound_executor.rs#L100-L115)
   
   #### A10. Task placement ignores where shuffle data lives
   **Impact:** Low–Medium · **Effort:** M
   
   Binding sorts executors by free vcores only, even though 
`partition_stats.num_bytes` for every input location is known at bind time, so 
most shuffle bytes cross the network unnecessarily on bandwidth-limited 
clusters.
   
   **Suggested change:** prefer the executor holding the largest share of a 
task's input bytes, with a short delay-scheduling window before falling back.
   
   **Where:** 
[`mod.rs:515-566`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L515-L566),
 
[`mod.rs:568-630`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L568-L630)
   
   ### Code
   
   #### C1. Sort-shuffle spill and final write block tokio worker threads
   **Impact:** High · **Effort:** S
   
   `spill_all_partitions` runs inline in the async drain loop (interleave, IPC 
encode, compress, file write), and `write_task_consolidated` copies every spill 
file and buffer into `data.arrow` inside the coordinator future. Both pin a 
runtime worker for the whole operation, stalling other partition pipelines on 
it.
   
   **Suggested change:** run the spill and the consolidated write in 
`spawn_blocking` (or the I/O pool from A9), moving the buffered state into the 
closure and back.
   
   **Where:** 
[`writer.rs:672-679`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/writer.rs#L672-L679),
 
[`writer.rs:1176-1186`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/writer.rs#L1176-L1186),
 
[`spill.rs:114-140`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/spill.rs#L114-L140)
   
   #### C15. A probe side read as broadcast runs the whole join in one task
   **Impact:** Medium · **Effort:** S · PR #2501
   
   In q08 stage 5, AQE produced a `CollectLeft` join whose probe child is a 
`ShuffleReaderExec` with `broadcast: true`. A broadcast reader reads all 
upstream partitions as one, so the stage ran as a single task pushing 80M probe 
rows through one hash join (7.4 s).
   
   **Suggested change:** set the broadcast flag only on the build child when 
join selection swaps sides or converts to `CollectLeft`, and add a plan-shape 
test.
   
   **Where:** #2419 `q08.md` stage 5, 
[`state/aqe/`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/state/aqe)
   
   #### C16. Date-partitioned scans list every partition directory
   **Impact:** Medium · **Effort:** S
   
   The fact tables are laid out as `lineitem/l_shipdate=YYYY-MM-DD/…`, but 
q06's scan (a one-year `l_shipdate` filter) starts its file list at 
`l_shipdate=1992-01-02`; q12 and q14 look the same. The filter only reaches 
row-group statistics, so every file in every date directory is listed, planned, 
opened and assigned a task.
   
   **Suggested change:** register the tables with their partition columns 
(`ListingOptions::with_table_partition_cols`, or the Iceberg provider) so 
directories are pruned before listing, and fix serialisation if Ballista drops 
partition columns. The harness that produced the #2419 layout isn't in 
`benchmarks/`, so where the fix lands needs confirming.
   
   **Where:** #2419 `q06.md`, `q12.md`, `q14.md`
   
   #### C2. Flight `do_get` decodes and re-compresses with hard-coded LZ4
   **Impact:** Medium · **Effort:** S
   
   The Flight path decodes every batch and re-encodes it with 
`CompressionType::LZ4_FRAME` regardless of 
`ballista.shuffle.compression.codec`. In the sort-shuffle branch, 
`MultiStreamPartitionStream` is polled directly on the gRPC server runtime with 
no `spawn_blocking`. Client result fetches always use this path; shuffle reads 
do when `remote_prefer_flight=true`.
   
   **Suggested change:** route client fetches to block transport; where Flight 
stays, use the `spawn_blocking` + channel pattern for sort-shuffle and the 
configured codec.
   
   **Where:** 
[`flight_service.rs:134-156`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/executor/src/flight_service.rs#L134-L156),
 
[`flight_service.rs:199-206`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/executor/src/flight_service.rs#L199-L206)
   
   #### C3. Client result fetch opens a new connection per partition with 64 
KiB windows
   **Impact:** Medium · **Effort:** S
   
   `fetch_partition` calls `BallistaClient::try_new(…, 0, 0)` per partition: a 
new TCP/HTTP2 (and TLS) handshake each time, and the default flow-control 
window caps throughput at ~64 KiB per RTT. Partitions are drained one at a 
time. This is the cheap part of A8.
   
   **Suggested change:** reuse one client per executor, pass the configured 
`initial_*_window_size`, and replace `flatten()` with `buffered(n)` (or 
`buffer_unordered(n)` when order doesn't matter).
   
   **Where:** 
[`distributed_query.rs:624-640`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/distributed_query.rs#L624-L640),
 
[`distributed_query.rs:881-892`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/distributed_query.rs#L881-L892)
   
   #### C4. Shuffle index and schema header are re-read on every partition fetch
   **Impact:** Medium · **Effort:** S
   
   Each partition request stats `data.arrow.index` and parses the whole index 
before serving bytes, including on the default block-transport path. Local 
reads and `do_get` also open the data file to decode the schema header, then 
open it again. With K reducers over M task files, each index is parsed K × M 
times.
   
   **Suggested change:** keep a small executor-side LRU of parsed 
`(ShuffleIndex, SchemaRef)` per data path, invalidated on job cleanup, and 
replace the `exists()` probe with an open that handles not-found (as the 
existing TODO says).
   
   **Where:** 
[`flight_service.rs:498-506`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/executor/src/flight_service.rs#L498-L506),
 
[`reader.rs:36-66`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/reader.rs#L36-L66),
 
[`shuffle_reader.rs:1267-1290`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L1267-L1290)
   
   #### C5. Passthrough shuffle writer writes small batches as-is
   **Impact:** Medium · **Effort:** S · PR #2518 (issue #2517)
   
   `write_stream_to_disk` writes each incoming batch as its own IPC message and 
compression frame, so small batches after selective filters or repartitions 
inflate file size, compression overhead and downstream decode calls.
   
   **Suggested change:** coalesce to `batch_size` before writing.
   
   **Where:** 
[`utils.rs:224-227`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/utils.rs#L224-L227),
 
[`shuffle_writer.rs:594-596`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_writer.rs#L594-L596)
   
   #### C6. Sort-shuffle end-of-input encode doubles peak memory
   **Impact:** Medium · **Effort:** M
   
   When a drain finishes, `encode_buffered_partitions` encodes all buffered 
rows into one `Vec<u8>` per partition while the source batches are still alive. 
The encoded buffers then stay resident, uncharged to the memory pool, until 
every other input finishes: a ~2× peak and invisible memory while waiting on 
the coordinator.
   
   **Suggested change:** write the final flush straight to a segment file like 
a spill (fits A6's segment index); otherwise drop each partition's source 
indices as it's encoded and grow the reservation by the encoded size.
   
   **Where:** 
[`writer.rs:40-80`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/writer.rs#L40-L80),
 
[`writer.rs:750-764`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/writer.rs#L750-L764)
   
   #### C7. Each revive clones the job map and write-locks every graph
   **Impact:** Medium · **Effort:** S
   
   `get_running_job_cache` builds a fresh `HashMap` of running jobs on every 
revive and poll, binders take each graph's write lock even with nothing 
pending, and `bind_one` re-walks the plan for `stage_has_input_collapse` and 
repeats the `BallistaConfig` lookup per task. This is the cheap subset of A4.
   
   **Suggested change:** check `available_tasks()` under a read lock (or an 
atomic counter) before taking the write lock, and compute `is_collapse` and the 
partition cap once when the `RunningStage` is created.
   
   **Where:** 
[`task_manager.rs:395-409`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/state/task_manager.rs#L395-L409),
 
[`mod.rs:433-444`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L433-L444),
 
[`mod.rs:538-545`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L538-L545),
 
[`mod.rs:591-598`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/scheduler/src/cluster/mod.rs#L591-L598)
   
   #### C8. Sort-shuffle local reads go through an unbuffered `File`
   **Impact:** Medium · **Effort:** S
   
   `MultiStreamPartitionStream` and the schema-header read build 
`StreamReader<File>` directly, so arrow's several small `read_exact` calls per 
message each become a syscall. This serves local sort-shuffle reads (about 1/N 
of input on an N-executor cluster) and `do_get` with 
`remote_prefer_flight=true`; the other shuffle formats already use a 256 KiB 
`BufReader`.
   
   **Suggested change:** wrap the file in a `BufReader` sized `min(256 KiB, 
segment length)` over `.take(end - start)` after the seek.
   
   **Where:** 
[`multi_stream_reader.rs:100-107`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/multi_stream_reader.rs#L100-L107),
 
[`reader.rs:66`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/reader.rs#L66)
   
   #### C9. Local shuffle files are read one at a time, with disk reads inside 
`poll_next`
   **Impact:** Low–Medium · **Effort:** S
   
   The local branch of `send_fetch_partitions` only opens files, serially, on 
one blocking thread; the read, decompress and decode happen later in 
`LocalShuffleStream::poll_next` on the task runtime, and streams are consumed 
in sequence through `try_flatten`.
   
   **Suggested change:** move the whole read to the blocking side with bounded 
concurrency (2–4 files) and send decoded batches through the channel.
   
   **Where:** 
[`shuffle_reader.rs:828-841`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L828-L841),
 
[`shuffle_reader.rs:583-592`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L583-L592),
 
[`shuffle_reader.rs:411-416`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L411-L416)
   
   #### C10. Remote fetch spawns one tokio task per location up front
   **Impact:** Low–Medium · **Effort:** S
   
   Every remote location gets a `SpawnedTask` immediately, most of which park 
on the three semaphores, and each acquisition locks a 
`std::sync::Mutex<HashMap>` to find the per-address semaphore. 10k upstream 
locations means 10k tasks and 10k lock acquisitions before any data moves.
   
   **Suggested change:** drive fetches from 
`stream::iter(locations).buffer_unordered(max_reqs)` with per-address 
sub-queues or a pre-built address→semaphore map.
   
   **Where:** 
[`shuffle_reader.rs:861-933`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/shuffle_reader.rs#L861-L933)
   
   #### C11. Every `PartitionLocation` carries full executor metadata
   **Impact:** Low–Medium · **Effort:** S
   
   Each location embeds a complete `ExecutorMetadata` (including 
`specification` and `os_info`) plus a `job_id` string. Sort shuffle produces M 
× K locations per stage, repeated across task plans, job status and scheduler 
memory, and adding to A3's encode cost.
   
   **Suggested change:** keep a per-plan executor table referenced by index (or 
send only id/host/port), and drop `job_id` from the per-location `PartitionId`.
   
   **Where:** 
[`ballista.proto:586-606`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/proto/ballista.proto#L586-L606),
 
[`ballista.proto:705-712`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/proto/ballista.proto#L705-L712)
   
   #### C17. Stage row metrics sum across operators and switch units
   **Impact:** Low · **Effort:** S
   
   Not a slowdown, but it makes the other findings hard to measure. Stage 
`output_rows` sums every operator's output (q08 s1 reports 11,999,979,418, 
exactly 2× its 5,999,989,709 input). Source-stage `input_rows` sometimes counts 
batches or groups (q01 s0 and q06 s0 report 401 and 256 for full lineitem 
scans). The event log has no per-stage task count and no per-task 
spill/write/fetch time.
   
   **Suggested change:** report the terminal operator's `output_rows` and leaf 
rows as stage input, and add task count and per-task `spill_time`, `write_time` 
and `fetch_time` percentiles to `JobEnd`.
   
   **Where:** #2419 `q01.md` s0, `q08.md` s1, `q17.md` s2, scheduler event log 
(#2264)
   
   #### C12. `ballista_config()` clones and re-parses settings on every call
   **Impact:** Low · **Effort:** S
   
   The accessor clones the whole settings `HashMap`, and every getter parses 
its string value with a fallback lookup into the defaults map. It's called per 
`execute()`, per input partition in writers and readers, and in `bind_one`.
   
   **Suggested change:** return `Arc<BallistaConfig>` (or a reference) and 
parse into typed fields once at construction.
   
   **Where:** 
[`extension.rs:484-490`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/extension.rs#L484-L490),
 
[`config.rs:899-935`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/config.rs#L899-L935)
   
   #### C13. View columns are garbage-collected unconditionally
   **Impact:** Low · **Effort:** S
   
   `compact_view_columns` calls `gc()` on every StringView/BinaryView column 
after each interleave, copying all data buffers even when the batch is already 
dense.
   
   **Suggested change:** compact only when the data buffers are much larger 
than the bytes the views reference (arrow's `BatchCoalescer` uses 2×).
   
   **Where:** 
[`partitioned_batch_iterator.rs:27-50`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/sort_shuffle/partitioned_batch_iterator.rs#L27-L50)
   
   #### C14. Pull-mode client polls job status every 50 ms
   **Impact:** Low · **Effort:** S
   
   With `ballista.client.pull=true` (push is the default), each short query 
gets ~25 ms average added latency, and every waiting client sends 20 RPCs/s to 
the scheduler.
   
   **Suggested change:** back off exponentially from ~5 ms to 250 ms, or remove 
pull mode now that push is the default.
   
   **Where:** 
[`distributed_query.rs:543-583`](https://github.com/apache/datafusion-ballista/blob/187e21fc/ballista/core/src/execution_plans/distributed_query.rs#L543-L583)
   
   **Describe alternatives you've considered**
   
   N/A. This is a tracking issue.
   
   **Additional context**
   
   Open questions from the SF1000 runs:
   
   - q05, q07, q18 and q21 got slower between the Sep 8 build (DataFusion 
55.0.0) and 32ceaf8 (55.1.0): +18%, +7%, +27% and +25% end to end. #2419's 
plans are from 54.0.0, so they can't be diffed against either build. Event-log 
plans from both builds for these four queries would show whether stage shapes 
changed or only per-stage time.
   - The lineitem writer tail (A6, C1, A5) has three candidate causes: 
spill-and-copy on the critical path, drains stalling behind a spill on a shared 
worker, or concurrent scan stages competing for disk and network. Per-task 
spill and write timings (C17) would separate them. A11 removes the stage from 
q08, q09 and q17 either way.
   - Rankings assume the default block transport for remote shuffle reads. On 
an N-executor cluster only ~1/N of shuffle locations are local, so local-read 
findings (C8, C9) rank lower here than they would on a single node.
   


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