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]

Reply via email to