You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
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<u8> 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.
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.rsandnative/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 aBatchCoalescer. In arrow 58.3, the normal path always doescopy_rowsinto in-progress builders, even when the pushed batch is exactlytarget_batch_sizeand the buffer is empty (the zero-copy bypass only activates whenwith_biggest_coalesce_batch_sizeis explicitly set, which we never do).PartitionedBatchIteratoralready emits exactlybatch_size-row batches frominterleave_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 andbatch.num_rows() >= batch_size, writing the batch directly (or configurebiggest_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 andconcat_batchesthem tobatch_size, then hands the result to the long-livedBufBatchWriter, whose ownBatchCoalescercopies 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 newzstd::Encoder(freshZSTD_CCtxplus workspace),lz4_flex::FrameEncoder, orsnap::FrameEncoderfor 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 sincewrite_batchtakes&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_newre-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 inShuffleBlockWriter::try_newand write blocks manually viaIpcDataGenerator::encoded_batch+write_message. This matters for wide schemas and many small blocks (e.g. 2000-partition shuffles) and also drops the per-blockStreamWriterallocation.Medium impact
Per-partition allocations in
PartitionedBatchIterator::new(partitioners/partitioned_batch_iterator.rs): for every non-empty partition it rebuilds aVec<&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 intoPartitionedBatchesProducerand convert indices perbatch_sizechunk into a reusable scratch buffer.BufBatchWriterchurn in multi-partition finish and spill:finish_partition(writers/local/local_partition_writer.rs) andSpillWriter::write(writers/local/spill.rs) construct a newBufBatchWriterper partition per event, so the internalVec<u8>regrows towardwrite_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 == 0means 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_buffersclonesArrayDataper 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_spillpushes indices in a manual loop afterreserve;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 branchlessr + ((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_writeafter input is exhausted, whileinsert_batchis 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.