NoahKusaba commented on PR #2518:
URL:
https://github.com/apache/datafusion-ballista/pull/2518#issuecomment-5939535515
Benchmark results for this change. I didn't add the benchmark to the PR; the
source is below if anyone wants to reproduce it.
**Setup:** `ShuffleWriterExec` writes 8 partitions × 131,072 rows (1M rows;
8 columns: 4 × Int64, 2 × Float64, 2 × Utf8), default LZ4 codec. Every case
slices one batch of distinct rows, so all cases write identical data and only
the input batch size changes. "write + read" also reads every file back with
`StreamReader`, the way a reducer does. Files are on tmpfs, so disk write-back
stays out of the timings. Times are the median of 3 rounds, with `main`
(187e21fc) and this branch run alternately.
| Input batches | Measure | main | this PR | Change |
|---|---|---|---|---|
| 128 rows | write | 74.1 ms | 41.8 ms | **−44%** |
| 128 rows | write + read | 132.7 ms | 73.3 ms | **−45%** |
| 8192 rows | write | 29.4 ms | 30.3 ms | +3% |
| 8192 rows | write + read | 58.2 ms | 58.0 ms | no change |
On disk, the 128-row input goes from 8,192 IPC batches (88.5 MB) to 128
(72.0 MB, **−19%**). With 8192-row input, both versions write the same files.
<details>
<summary><code>benchmarks/benches/passthrough_shuffle.rs</code></summary>
```rust
//! Criterion benchmarks for the passthrough shuffle writer.
//!
//! The passthrough writer (`ShuffleWriterExec`) writes each child partition
//! to its own Arrow IPC file. These benchmarks feed it the same rows either
//! as many small batches (as produced after a selective filter) or as
//! `batch_size` batches, and measure the write alone and the write followed
//! by reading every file back, as a reducer would.
use std::fs::File;
use std::io::BufReader;
use std::sync::Arc;
use ballista_core::execution_plans::ShuffleWriterExec;
use criterion::{Criterion, criterion_group, criterion_main};
use datafusion::arrow::array::{Array, Float64Array, Int64Array, StringArray};
use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use datafusion::arrow::ipc::reader::StreamReader;
use datafusion::arrow::record_batch::RecordBatch;
use datafusion::datasource::memory::MemorySourceConfig;
use datafusion::datasource::source::DataSourceExec;
use datafusion::physical_plan::ExecutionPlan;
use datafusion::prelude::SessionContext;
use futures::TryStreamExt;
use rand::RngExt;
use tempfile::TempDir;
use tokio::runtime::Runtime;
const NUM_PARTITIONS: usize = 8;
const ROWS_PER_PARTITION: usize = 131_072;
fn build_schema() -> SchemaRef {
Arc::new(Schema::new(vec![
Field::new("i0", DataType::Int64, false),
Field::new("i1", DataType::Int64, false),
Field::new("i2", DataType::Int64, false),
Field::new("i3", DataType::Int64, false),
Field::new("f0", DataType::Float64, false),
Field::new("f1", DataType::Float64, false),
Field::new("s0", DataType::Utf8, false),
Field::new("s1", DataType::Utf8, false),
]))
}
fn build_batch(schema: &SchemaRef, rows: usize) -> 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..rows).map(|_| rng.random()).collect();
columns.push(Arc::new(Int64Array::from(vals)));
}
for _ in 0..2 {
let vals: Vec<f64> = (0..rows).map(|_| rng.random()).collect();
columns.push(Arc::new(Float64Array::from(vals)));
}
for _ in 0..2 {
let vals: Vec<String> = (0..rows)
.map(|_| format!("value-{}", rng.random_range(0..1_000_000)))
.collect();
columns.push(Arc::new(StringArray::from(vals)));
}
RecordBatch::try_new(schema.clone(), columns).unwrap()
}
/// `NUM_PARTITIONS` partitions of the same `ROWS_PER_PARTITION` distinct
rows,
/// delivered as slices of `rows_per_batch` rows, so every input size writes
/// identical data.
fn create_input(rows: &RecordBatch, rows_per_batch: usize) -> Arc<dyn
ExecutionPlan> {
let schema = rows.schema();
let partition: Vec<RecordBatch> = (0..ROWS_PER_PARTITION)
.step_by(rows_per_batch)
.map(|offset| rows.slice(offset, rows_per_batch))
.collect();
let partitions = vec![partition; NUM_PARTITIONS];
let memory_source =
Arc::new(MemorySourceConfig::try_new(&partitions, schema,
None).unwrap());
Arc::new(DataSourceExec::new(memory_source))
}
/// Runs the writer over every output partition and returns the file paths.
fn run_write(rt: &Runtime, input: Arc<dyn ExecutionPlan>, work_dir: &str) ->
Vec<String> {
let writer = Arc::new(
ShuffleWriterExec::try_new("bench_job".into(), 1, input,
work_dir.to_string())
.unwrap(),
);
let task_ctx = SessionContext::new().task_ctx();
rt.block_on(async {
let streams = (0..NUM_PARTITIONS)
.map(|p| writer.execute(p, task_ctx.clone()).unwrap())
.map(|s| s.try_collect::<Vec<_>>());
let summaries =
futures::future::try_join_all(streams).await.unwrap();
summaries
.iter()
.flatten()
.flat_map(|b| {
let paths =
b.column(1).as_any().downcast_ref::<StringArray>().unwrap();
paths
.iter()
.map(|p| p.unwrap().to_string())
.collect::<Vec<_>>()
})
.collect()
})
}
/// Reads every file back and returns `(rows, ipc_batches)`.
fn read_back(paths: &[String]) -> (usize, usize) {
let mut rows = 0;
let mut batches = 0;
for path in paths {
let file = BufReader::new(File::open(path).unwrap());
for batch in StreamReader::try_new(file, None).unwrap() {
rows += batch.unwrap().num_rows();
batches += 1;
}
}
(rows, batches)
}
fn bench_passthrough_shuffle(c: &mut Criterion) {
let rt = Runtime::new().unwrap();
let work_dir = TempDir::new().unwrap();
let dir = work_dir.path().to_str().unwrap();
let rows = build_batch(&build_schema(), ROWS_PER_PARTITION);
for (name, rows_per_batch) in [("128_row_batches", 128),
("8192_row_batches", 8192)] {
let input = create_input(&rows, rows_per_batch);
let paths = run_write(&rt, input.clone(), dir);
let bytes: u64 = paths
.iter()
.map(|p| std::fs::metadata(p).unwrap().len())
.sum();
let (rows, batches) = read_back(&paths);
println!(
"{name}: {rows} rows, {batches} IPC batches, {bytes} bytes on
disk across {} files",
paths.len()
);
let mut group = c.benchmark_group("passthrough_shuffle");
group.sample_size(20);
group.bench_function(format!("write/{name}"), |b| {
b.iter(|| run_write(&rt, input.clone(), dir));
});
group.bench_function(format!("write_and_read/{name}"), |b| {
b.iter(|| read_back(&run_write(&rt, input.clone(), dir)));
});
group.finish();
}
}
criterion_group!(benches, bench_passthrough_shuffle);
criterion_main!(benches);
```
Register it in `benchmarks/Cargo.toml` with `[[bench]] harness = false, name
= "passthrough_shuffle"`, then run `TMPDIR=/dev/shm cargo bench -p
ballista-benchmarks --bench passthrough_shuffle`.
</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]