apache / apache/datafusion-comet

perf: optimize native shuffle writer (redundant copies, per-block allocations)

Open
#5,002 2 comments 2 reactions 1 assignee Claimed by @andygrove View on GitHub
enhancement performance
Dominant language
Scala
Stars
1.3k
Forks
373
Avg merge
2d 4h
Merged PRs (30d)
198

Description

### What is the problem the feature request solves?

A review of the native shuffle writer (`native/shuffle/`) identified several performance optimization opportunities. The partitioning and interleave logic is in good shape (reused scratch buffers, counting-sort partition assignment, `interleave_record_batch`, dedup'd buffer accounting), but the serialization/write side has redundant copies, per-block allocations, and avoidable per-partition churn.

The crate already has benchmarks (`native/shuffle/benches/shuffle_writer.rs` and `native/shuffle/src/bin/shuffle_bench.rs`), so each item below can be validated in isolation before running Spark-level benchmarks.

### High impact

- [ ] **Every batch is copied an extra time through `BatchCoalescer`, even when already full.** `BufBatchWriter::write` (`writers/buf_batch_writer.rs`) pushes every batch through a `BatchCoalescer`. In arrow 58.3, the normal path always does `copy_rows` into in-progress builders, even when the pushed batch is exactly `target_batch_size` and the buffer is empty (the zero-copy bypass only activates when `with_biggest_coalesce_batch_size` is explicitly set, which we never do). `PartitionedBatchIterator` already emits exactly `batch_size`-row batches from `interleave_record_batch` (except the tail), so in the multi-partition spill and finish paths essentially all shuffle data is copied twice: once by interleave, once by the coalescer. Fix: bypass the coalescer when its buffer is empty and `batch.num_rows() >= batch_size`, writing the batch directly (or configure `biggest_coalesce_batch_size`). This removes one full copy of the entire shuffle payload with bit-identical output.

- [ ] **The single-partition path coalesces twice.** `SinglePartitionShufflePartitioner` (`partitioners/single_partition.rs`) buffers small batches and `concat_batches` them to `batch_size`, then hands the result to the long-lived `BufBatchWriter`, whose own `BatchCoalescer` copies all the rows again. The concat layer appears to predate the coalescer in the writer and is now redundant. Dropping it removes another full copy in single-partition shuffles.

- [ ] **A fresh compression context is allocated per IPC block.** `ShuffleBlockWriter::write_batch` (`writers/shuffle_block_writer.rs`) creates a new `zstd::Encoder` (fresh `ZSTD_CCtx` plus workspace), `lz4_flex::FrameEncoder`, or `snap::FrameEncoder` for every block. With default 8192-row blocks, and especially with high partition counts where blocks are small, context setup is a meaningful fraction of compression time. zstd explicitly documents context reuse as a major win. Requires a small refactor since `write_batch` takes `&self`. The read side (`ipc.rs`) has the same per-block decoder allocation.

- [ ] **The IPC schema is re-encoded for every block.** Each block is a standalone IPC stream, so `StreamWriter::try_new` re-serializes the schema flatbuffer per block. The bytes must be in the stream, but the encoding work does not need repeating: pre-encode the schema message once in `ShuffleBlockWriter::try_new` and write blocks manually via `IpcDataGenerator::encoded_batch` + `write_message`. This matters for wide schemas and many small blocks (e.g. 2000-partition shuffles) and also drops the per-block `StreamWriter` allocation.

### Medium impact

- [ ] **Per-partition allocations in `PartitionedBatchIterator::new`** (`partitioners/partitioned_batch_iterator.rs`): for every non-empty partition it rebuilds a `Vec<&RecordBatch>` over all buffered batches, and up-converts the partition's entire `(u32, u32)` index list into a fresh `(usize, usize)` Vec, re-materializing 16 bytes per row of index data per write cycle. Hoist the batch-ref vec into `PartitionedBatchesProducer` and convert indices per `batch_size` chunk into a reusable scratch buffer.

- [ ] **`BufBatchWriter` churn in multi-partition finish and spill**: `finish_partition` (`writers/local/local_partition_writer.rs`) and `SpillWriter::write` (`writers/local/spill.rs`) construct a new `BufBatchWriter` per partition per event, so the internal `Vec` regrows toward `write_buffer_size` (1MB default) each time and is dropped. With 2000 partitions that is 2000 buffer growth cycles per pass. Reuse a buffer across partitions.

- [ ] **RoundRobin hashes all columns by default** (`partitioners/multi_partition.rs`): `max_hash_columns == 0` means every column of every row is murmur3-hashed just to get an even, deterministic spread. Spark's round robin does no content hashing at all; distribution quality only needs a little entropy. A small default cap (1 or 2 columns) would cut CPU substantially on wide rows while keeping determinism. The config knob already exists. Caveat: changes partition assignments, which round robin is allowed to do.

- [ ] **`count_new_buffers` clones `ArrayData` per column per batch** (`partitioners/multi_partition.rs`): `column.to_data()` recursively clones the ArrayData tree (Vec allocations plus Arc bumps) on the insert hot path. Cheap relative to hashing, but avoidable for the common flat-schema case.

### Micro

- [ ] `buffer_partitioned_batch_may_spill` pushes indices in a manual loop after `reserve`; `indices.extend(row_indices.iter().map(|&r| (buffered_partition_idx, r)))` gives the trusted-len specialization.
- [ ] `pmod` (`comet_partitioning.rs`) does up to two `%` ops per row; a branchless `r + ((r >> 31) & n)` on the i32 remainder saves one. The hash and pmod passes over the batch could also be fused.

### Structural observation (not a defect)

In the no-spill case all interleave, compression, and IO happen serially inside `shuffle_write` after input is exhausted, while `insert_batch` is nearly pure buffering. This is a deliberate design (it maximizes coalescing per partition), but it means shuffle write time is fully exposed on the critical path rather than overlapped with upstream compute.

### Describe the potential solution

Address the checklist items above in separate PRs, starting with the two copy-elimination items (bit-identical output, safest wins), then compression context reuse, which is likely the biggest single CPU win for zstd users.

### Additional context

Found during a code review of the shuffle write path: `ShuffleWriterExec` -> `MultiPartitionShuffleRepartitioner` / `SinglePartitionShufflePartitioner` -> `PartitionedBatchIterator` -> `BufBatchWriter` / `ShuffleBlockWriter` -> `LocalPartitionWriter` / `SpillWriter`.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.