NoahKusaba opened a new pull request, #2518: URL: https://github.com/apache/datafusion-ballista/pull/2518
## Which issue does this PR close? - Closes #2517. ## Rationale for this change The passthrough shuffle writer (`ShuffleWriterExec`) writes every batch it receives as its own Arrow IPC message, with its own compression frame. After a selective filter or a `RepartitionExec`, these batches can be a few hundred rows, which makes shuffle files bigger, adds a compression frame per tiny batch, and gives downstream readers more batches to decode. ## What changes are included in this PR? - `utils::write_stream_to_disk` runs batches through arrow's `BatchCoalescer` before `writer.write`, targeting the session `batch_size` (`datafusion.execution.batch_size`, default 8192). This happens inside the existing `spawn_blocking` writer task, so the tokio workers do no extra work. - The coalescer is built with `with_biggest_coalesce_batch_size(Some(batch_size / 2))`, so a batch of more than half the target size is written as-is when nothing is buffered. DataFusion uses the same setting in `LimitedBatchCoalescer::new` ([`coalesce/mod.rs#L63-L64`](https://github.com/apache/datafusion/blob/7d3835c71f30cbd3c3ae4041732267f1f453097a/datafusion/physical-plan/src/coalesce/mod.rs#L63-L64)), which `CoalesceBatchesExec` uses ([`coalesce_batches.rs#L230`](https://github.com/apache/datafusion/blob/7d3835c71f30cbd3c3ae4041732267f1f453097a/datafusion/physical-plan/src/coalesce_batches.rs#L230)). Links are to DataFusion `55.1.0`. - `PartitionStats.num_batches` now counts the batches actually written to the file, not the batches received. - The `sort_shuffle` module doc is corrected: output is one `data.arrow` + `data.arrow.index` pair per task (not per input partition), and every hash-repartitioning stage goes through it. Rows can now be held in the coalescer, up to `batch_size` rows per output partition, until the stream ends. The bounded channel feeding the writer already held batches the same way. ## Are there any user-facing changes? No config changes. Passthrough shuffle files contain fewer, larger batches, and the reported `num_batches` drops to match. **API change:** the public `ballista_core::utils::write_stream_to_disk` takes a new `batch_size: usize` argument. ## What is the testing strategy for this PR? - New `utils::tests::write_stream_to_disk_coalesces_small_batches`: writes 10 batches of 3 rows with `batch_size = 8`, reads the file back and checks for exactly `[0..8, 8..16, 16..24, 24..30]`, plus the stats (30 rows, 4 batches, bytes equal to the file size). - New `utils::tests::write_stream_to_disk_passes_large_batches_through`: two 20-row batches with `batch_size = 8` come back unchanged as 2 batches. - Updated `shuffle_reader::tests::test_read_local_shuffle`: it previously expected the writer's two 3-row input batches back separately, and now expects the single coalesced 6-row batch. - `cargo test` passes for `ballista-core`, `ballista-executor` and `ballista-scheduler`; `cargo clippy --all-targets -D warnings` and `cargo fmt --check` are clean. -- 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]
