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]
