NoahKusaba opened a new pull request, #2522: URL: https://github.com/apache/datafusion-ballista/pull/2522
## Which issue does this PR close? - Closes #2521. ## Rationale for this change `compact_view_columns` runs `gc()` on every `Utf8View` / `BinaryView` column after each interleave in the sort-shuffle writer. `gc()` copies every referenced byte, even when the column is already dense, or has no data buffers at all because every value is inlined (12 bytes or fewer). ## What changes are included in this PR? - `compact_view_columns` now compacts a column only when its data buffers are more than twice the bytes its views reference. Arrow's `BatchCoalescer` uses the same threshold before copying strings ([`byte_view.rs#L376-L379`](https://github.com/apache/arrow-rs/blob/782e5a685501a9db6cc8e9a3b7cbff894940c47a/arrow-select/src/coalesce/byte_view.rs#L376-L379), arrow 59.2.0). Two differences from arrow's version: - It sums buffer `len()` rather than `capacity()`, because `len()` is what the IPC writer serializes. - A column with data buffers but no referenced bytes (all values inlined) is still compacted, which drops the unused buffers. - Columns with no data buffers are never compacted. **Tradeoff:** a column that skips compaction can carry up to 2× its referenced string bytes into the shuffle file. In the dense benchmark case below, files are 1.9% larger. ## Are there any user-facing changes? No. ## What is the testing strategy for this PR? - New `dense_view_columns_are_not_copied`: a dense `Utf8View` batch comes back as the same array (`Arc::ptr_eq`) and equal to the input. - New `sparse_view_columns_are_compacted`: picking 1 row out of 100 long strings returns exactly that row, backed by a single data buffer of exactly that string's length. - `cargo test -p ballista-core`, clippy and fmt pass. Note: CI's Clippy job is failing on all PRs right now because of Rust 1.99. That's fixed separately in #2519. ### Benchmark `SortShuffleWriterExec` over 50 batches × 8192 rows (4 × Int64, 4 × Utf8View; distinct rows in every batch), files on tmpfs. Times are the median of 3 rounds, with `main` (187e21fc) and this branch run alternately. | Case | main | this PR | Time | On disk | |---|---|---|---|---| | Inline strings (≤ 12 bytes), 200 output partitions | 189.3 ms | 171.8 ms | **−9%** | same | | Long strings, 200 output partitions | 325.3 ms | 320.4 ms | within noise | same | | Long strings, 1 output partition | 209.9 ms | 200.1 ms | **−5%** | +1.9% | - **Long strings, 200 partitions:** each output batch takes about 1/200 of the rows from each input batch, so its columns are sparse and both versions compact them. - **Inline strings and the 1-partition case:** the columns are already dense or bufferless, so the copy is skipped. I didn't add the benchmark to this PR. The source is below for anyone who wants to reproduce it. <details> <summary><code>benchmarks/benches/sort_shuffle_view_gc.rs</code></summary> ```rust //! Criterion benchmarks for the sort-shuffle writer over `Utf8View` columns. //! //! The writer interleaves rows into per-partition batches and compacts view //! columns afterwards. These cases cover columns with no data buffers //! (all strings inlined), long strings spread over many output partitions //! (sparse after interleave) and long strings into one output partition //! (dense after interleave). use std::sync::Arc; use ballista_core::execution_plans::SortShuffleWriterExec; use ballista_core::execution_plans::sort_shuffle::SortShuffleConfig; use criterion::{Criterion, criterion_group, criterion_main}; use datafusion::arrow::array::{Array, Int64Array, StringViewArray}; use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; use datafusion::arrow::record_batch::RecordBatch; use datafusion::datasource::memory::MemorySourceConfig; use datafusion::datasource::source::DataSourceExec; use datafusion::physical_plan::coalesce_partitions::CoalescePartitionsExec; use datafusion::physical_plan::expressions::Column; use datafusion::physical_plan::{ExecutionPlan, Partitioning}; use datafusion::prelude::SessionContext; use futures::TryStreamExt; use rand::RngExt; use tempfile::TempDir; use tokio::runtime::Runtime; const BATCH_SIZE: usize = 8192; const NUM_BATCHES: usize = 50; fn build_schema() -> SchemaRef { let mut fields: Vec<Field> = (0..4) .map(|i| Field::new(format!("i{i}"), DataType::Int64, false)) .collect(); fields.extend((0..4).map(|i| Field::new(format!("s{i}"), DataType::Utf8View, false))); Arc::new(Schema::new(fields)) } /// One batch of distinct rows. `long` strings are 32+ bytes, so their bytes /// live in data buffers; short ones are at most 12 bytes and stay inline. fn build_batch(schema: &SchemaRef, long: bool) -> RecordBatch { let mut rng = rand::rng(); let mut columns: Vec<Arc<dyn Array>> = Vec::with_capacity(8); for _ in 0..4 { let vals: Vec<i64> = (0..BATCH_SIZE).map(|_| rng.random()).collect(); columns.push(Arc::new(Int64Array::from(vals))); } for _ in 0..4 { let vals = (0..BATCH_SIZE).map(|_| { let n = rng.random_range(0..10_000_000u32); if long { format!("customer-comment-{n:010}-padding") } else { format!("c{n}") } }); columns.push(Arc::new(StringViewArray::from_iter_values(vals))); } RecordBatch::try_new(schema.clone(), columns).unwrap() } fn create_input(long: bool) -> Arc<dyn ExecutionPlan> { let schema = build_schema(); let partition: Vec<RecordBatch> = (0..NUM_BATCHES) .map(|_| build_batch(&schema, long)) .collect(); let memory_source = Arc::new(MemorySourceConfig::try_new(&[partition], schema, None).unwrap()); Arc::new(CoalescePartitionsExec::new(Arc::new(DataSourceExec::new( memory_source, )))) } fn run_sort_shuffle( rt: &Runtime, input: Arc<dyn ExecutionPlan>, work_dir: &str, num_output_partitions: usize, ) { let writer = SortShuffleWriterExec::try_new( "bench_job".into(), 1, input, work_dir.to_string(), Partitioning::Hash(vec![Arc::new(Column::new("i0", 0))], num_output_partitions), SortShuffleConfig::new(true, BATCH_SIZE), ) .unwrap(); let task_ctx = SessionContext::new().task_ctx(); rt.block_on(async { let mut stream = writer.execute(0, task_ctx).unwrap(); while let Some(_batch) = stream.try_next().await.unwrap() {} }); } fn dir_bytes(dir: &std::path::Path) -> u64 { std::fs::read_dir(dir) .unwrap() .map(|e| { let e = e.unwrap(); if e.file_type().unwrap().is_dir() { dir_bytes(&e.path()) } else { e.metadata().unwrap().len() } }) .sum() } fn bench_view_gc(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("sort_shuffle_view_gc"); group.sample_size(20); for (name, long, num_output_partitions) in [ ("inline_strings_200_partitions", false, 200), ("long_strings_200_partitions", true, 200), ("long_strings_1_partition", true, 1), ] { let input = create_input(long); let work_dir = TempDir::new().unwrap(); let dir = work_dir.path().to_str().unwrap(); run_sort_shuffle(&rt, input.clone(), dir, num_output_partitions); println!("{name}: {} bytes on disk", dir_bytes(work_dir.path())); group.bench_function(name, |b| { b.iter(|| run_sort_shuffle(&rt, input.clone(), dir, num_output_partitions)); }); } group.finish(); } criterion_group!(benches, bench_view_gc); criterion_main!(benches); ``` Register it in `benchmarks/Cargo.toml` with `[[bench]] harness = false, name = "sort_shuffle_view_gc"`, then run `TMPDIR=/dev/shm cargo bench -p ballista-benchmarks --bench sort_shuffle_view_gc`. </details> -- 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]
