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]

Reply via email to