perf(core): write one sort-shuffle file per task instead of one per input - #2316
Conversation
71269a0 to
b6bd788
Compare
…nput
A task that owns several input partitions wrote one `data.arrow` plus index
per input partition. Every downstream reader fetching partition k therefore
opened one file per input partition, and a stage left M files behind where M
is the stage's input partition count.
The writer now emits a single file per task. Each input partition still
buckets and spills concurrently, and — importantly — still encodes its own
buckets to IPC bytes on its own task, so the interleave, framing and
compression stay parallel. The coordinator, which already awaited every input
before responding, then concatenates the finished buffers into one file and
writes one index.
The file is laid out partition-major: the schema header, then output
partition 0's bytes from every input in turn, then partition 1's, and so on.
Keeping each output partition contiguous is what lets the index stay one
offset per partition, so the reader is unchanged — `create_shuffle_path`
already resolves a sort-shuffle summary to `{stage_id}/{file_id}/data.arrow`,
and `MultiStreamPartitionStream` already crosses concatenated IPC streams
inside a byte range.
An earlier revision of this change did the encoding in the coordinator, which
collapsed write parallelism from P to 1 and cost 29% on TPC-H SF10. Encoding
per input is what makes the file-count reduction free.
Spill directories move from `{stage}/{file_id}/spill` to
`{stage}/{task_id}/spill-{input}`, so a task owns exactly one directory under
the stage and cleanup no longer strands an empty directory per input.
TPC-H SF10, 2 executors x 4 vcores, `--partitions 16`,
`max_partitions_per_task=0`. Old and new binaries alternate within each round
so drift is shared, 8 runs each:
| | median | sort-shuffle files |
|---|---|---|
| one file per input partition | 21.37 s | 1762 |
| one file per task | 20.05 s | 442 |
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
b6bd788 to
52a511f
Compare
|
Benchmarked this at SF1000 on EKS. Short version: it's a solid win, and I think the mechanism is more interesting than the headline number. SetupTPC-H SF1000 (ZSTD parquet, 32 MiB row groups) on S3, 32 executor pods x 8 vcores = 256 vcores across 4x r6i.24xlarge, 48 GiB memory pool per pod, Baseline is main at this PR's merge-base, so the only difference between the two builds is the two files in this PR. I crossed that with Seven queries: the four heaviest sort-shuffle writers, plus Q1 and Q6 as controls. Q6 plans Wall clock, seconds
At the shipping default the PR is 8.0% faster overall, with Q3 -22.3%, Q10 -22.8%, Q21 -13.7%. Spill volumes come out byte-identical between the two builds at each setting (Q9: 416.9 GB over 1790 events either way), which is a nice confirmation that this changes file layout only, not bucketing or spill decisions. Hypothesis: the win is parallel encoding, not throughputWhat the gather + IPC framing + lz4 work has to cross is the end-of-task serial path. There look to be three ways it can avoid that:
That predicts the PR should help most exactly where spilling wasn't already helping, and the numbers line up:
The part I find most useful is what this does to the spill budget as a tuning knob. On main it's a real trade: turning it on gains 25% on Q9 but costs 15% on Q3 and 12% on Q10. With this PR the setting stops mattering much (totals 239.58 vs 237.73, under 1%). Operators currently have to pick a side; after this they mostly don't. That reads to me as a better argument for the change than the 8%. Two caveats on reading the aboveOne iteration per cell. Q6 ranged 7.31-8.45s across the five runs, so the noise floor is around 8%. The double-digit per-query deltas clear that comfortably; the 8.0% total does not, and neither Q9's +3.0% nor Q18's -1.6% should be read as anything. Happy to rerun with 3 iterations if that's worth having. Also worth flagging because I nearly quoted it: the |
avantgardnerio
left a comment
There was a problem hiding this comment.
Nice work!
I originally tried to encode to a single file when implementing MPTs, but it ran much slower because I had encoding on the serial path. This PR keeps the encoding parallel, regardless of spill, and that's a big win (see perf numbers below).
I gave it a thorough review, and the only nit is that it would be worth revisiting the metrics at some point, but that doesn't have to be in this PR.
Thanks for the review! 🚀 |
Summary
A task that owns several input partitions wrote one
data.arrowplus indexper input partition. Every downstream reader fetching partition
ktherefore opened one file per input partition, and a stage left
Mfilesbehind, where
Mis the stage's input partition count.The writer now emits one file per task.
This matters more since
ballista.scheduler.max_partitions_per_taskdefaultsto
0(#2315): a task now routinely owns several input partitions, soMfiles per stage became
M/P.Design
Each input partition still buckets and spills concurrently, and still encodes
its own buckets to IPC bytes on its own task, so the interleave, framing
and compression stay parallel across inputs. The coordinator — which already
awaited every input before responding — then concatenates the finished buffers
into one file and writes one index.
The file is laid out partition-major: the schema header, then output
partition 0's bytes from every input in turn, then partition 1's, and so on.
Keeping each output partition contiguous is what lets the index stay one
offset per partition, which is why the reader is unchanged:
create_shuffle_pathalready resolves a sort-shuffle summary to{stage_id}/{file_id}/data.arrow, andMultiStreamPartitionStreamalreadycrosses concatenated IPC streams inside a byte range.
Spill directories move from
{stage}/{file_id}/spillto{stage}/{task_id}/spill-{input}, so a task owns exactly one directory underthe stage and cleanup no longer strands an empty directory per input partition.