peterxcli commented on PR #6440:
URL: 
https://github.com/apache/datafusion-comet/pull/6440#issuecomment-5903927651

   Benchmarked base `4594e2b027262a9b1e8942be0ab1ceb99133be98` against PR head 
`e3a21a8762d86327eb79a12ede10dc64cba5b406` using the existing standalone native 
shuffle benchmark.
   
   Median elapsed time and process peak RSS:
   
   | Output partitions / concurrent tasks | Timed runs per revision | Write 
seconds, base → PR | Time change | Peak RSS MiB, base → PR |
   | --- | ---: | ---: | ---: | ---: |
   | 200 / 1 | 5 | 0.204 → 0.208 | +2.0% | 30.0 → 29.6 |
   | 16,000 / 1 | 5 | 3.820 → 3.581 | -6.3% | 92.0 → 91.1 |
   | 16,000 / 8 | 3 | 55.412 → 55.738 | +0.6% | 412.8 → 398.5 |
   
   Allocated capacity of the inner spill-range vectors, measured before and 
after merging all partitions:
   
   | Output partitions / concurrent tasks | Base MiB, before → after | PR MiB, 
before → after | Outer table MiB, both revisions |
   | --- | ---: | ---: | ---: |
   | 200 / 1 | 0.391 → 0.391 | 0.391 → 0 | 0.005 |
   | 16,000 / 1 | 31.246 → 31.246 | 31.246 → 0 | 0.366 |
   | 16,000 / 8 | 249.969 → 249.969 | 249.969 → 0 | 2.930 |
   
   The PR releases all inner range capacity after merge in every run. The outer 
table remains until the writer is dropped. For eight tasks these are sums of 
each task's snapshots, not simultaneous process memory measurements. Peak RSS 
and elapsed time are separate measurements; this experiment does not establish 
executor memory savings or a general speedup.
   
   Each task processed 2,883,584 rows through 88 forced spills. Every output 
was decoded outside the timed process and checked for missing/duplicate rows, 
payload values, and valid partition boundaries. Data files were byte-identical 
between revisions and repeats for each workload.
   
   <details>
   <summary>Method, raw samples, and reproduction</summary>
   
   Ranges across timed runs:
   
   | Output partitions / concurrent tasks | Base seconds, min–max | PR seconds, 
min–max | Base peak RSS MiB, min–max | PR peak RSS MiB, min–max |
   | --- | ---: | ---: | ---: | ---: |
   | 200 / 1 | 0.191–0.212 | 0.188–0.233 | 29.2–31.7 | 29.2–31.8 |
   | 16,000 / 1 | 3.728–4.772 | 3.440–4.756 | 90.7–93.2 | 85.2–92.4 |
   | 16,000 / 8 | 52.049–56.132 | 51.212–62.168 | 401.7–419.0 | 393.2–475.6 |
   
   - Apple M4, 10 logical CPUs, 24 GiB RAM; macOS 27.0 (26A428); Rust 1.96.0, 
aarch64-apple-darwin. Local laptop, not an isolated benchmark host.
   - Release build with thin LTO and one codegen unit. Same benchmark 
instrumentation applied to both revisions in temporary worktrees; production PR 
code is unchanged.
   - One untimed warmup per revision/workload, then fresh processes in 
alternating base/PR order. The warmups are included in the CSV but excluded 
from medians.
   - One incomplete trial was stopped by a session interruption and rerun; only 
completed, validated trials are included.
   - One scan partition; 32,768 rows per batch; hash partitioning on column 0; 
LZ4; 1 MiB write buffer; `--max-buffer-bytes 1` forces 88 spill rounds per 
task. Eight tasks each process the full fixture.
   - Each task uses its own DataFusion runtime with the default memory pool. 
This does not exercise Spark/JNI pool overhead, shared executor fairness, or 
adaptive spill decisions under a memory limit.
   - Write elapsed time is reported by `shuffle_bench`; `/usr/bin/time -l` 
measures maximum process RSS. Instrumented merge time spans the first through 
last `finish_partition`, excluding `finish_all`'s final flush. Concurrent merge 
CSV values are per-run medians over eight tasks, not independent trials.
   - The instrumentation calculates inner `Vec` capacity directly in bytes; it 
does not equate allocator capacity with RSS. Readback and SHA256 validation run 
after the timed process.
   
   Fixture generated with DuckDB 1.4.3:
   
   ```sql
   COPY (
     SELECT i::BIGINT AS id, (i * 17 + 11)::BIGINT AS value
     FROM range(2883584) t(i)
   ) TO '/tmp/comet-6440-bench/input.parquet'
   (FORMAT PARQUET, COMPRESSION SNAPPY, ROW_GROUP_SIZE 131072);
   ```
   
   Fixture SHA256: 
`06e56434485b56b93e531c6242f636f40e96d87fa815a49d052fa2f6f15f1271`.
   
   Apply the patch below to both revisions and build each in its own target 
directory:
   
   ```sh
   cd native
   cargo build --release -p datafusion-comet-shuffle --features shuffle-bench 
--bin shuffle_bench --offline
   COMET_BENCH_RETAIN_OUTPUT=1 TOKIO_WORKER_THREADS=10 /usr/bin/time -l \
     target/release/shuffle_bench --input /tmp/comet-6440-bench/input.parquet \
     --batch-size 32768 --partitions 16000 --codec lz4 --hash-columns 0 \
     --max-buffer-bytes 1 --iterations 1 --write-buffer-size 1048576 \
     --concurrent-tasks 8 --output-dir /tmp/comet-6440-bench/repro-output
   ```
   
   Use 200 or 16,000 partitions and 1 or 8 concurrent tasks as listed above. 
Decode each retained `data.out` separately with `shuffle_bench --input ... 
--validate-output ...`; use a fresh output directory for each run.
   
   Raw samples (warmup rows excluded from summary):
   
   ```csv
   
case,variant,repeat,warmup,write_secs,max_rss_mib,merge_task_median_secs,inner_before_mib_sum,inner_after_mib_sum,outer_mib_sum,spills_per_task,rows_per_task,output_bytes_per_task,output_sha256
   
200p_1task,base,warmup,True,0.199,29.265625,0.018479,0.390625,0.390625,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,pr,warmup,True,0.209,29.328125,0.028284,0.390625,0.0,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,base,1,False,0.191,29.75,0.017124,0.390625,0.390625,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,pr,1,False,0.188,29.25,0.015261,0.390625,0.0,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,pr,2,False,0.195,29.28125,0.022733,0.390625,0.0,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,base,2,False,0.191,29.25,0.019934,0.390625,0.390625,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,base,3,False,0.211,31.65625,0.021414,0.390625,0.390625,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,pr,3,False,0.208,31.828125,0.031372,0.390625,0.0,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,pr,4,False,0.233,29.765625,0.02725,0.390625,0.0,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,base,4,False,0.204,30.0,0.025772,0.390625,0.390625,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,base,5,False,0.212,30.5625,0.020558,0.390625,0.390625,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
200p_1task,pr,5,False,0.211,29.578125,0.023424,0.390625,0.0,0.00457763671875,88,2883584,33323560,3c4cfe3c9f75546772804bd175ee6322c08e3ba0220a7201ba40ffbd93f7ed81
   
16000p_1task,base,warmup,True,3.948,86.140625,0.684438,31.24609375,31.24609375,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,pr,warmup,True,3.707,90.8125,0.707806,31.24609375,0.0,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,base,1,False,4.772,92.109375,0.609952,31.24609375,31.24609375,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,pr,1,False,4.756,91.484375,1.17317,31.24609375,0.0,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,pr,2,False,4.19,85.21875,0.590631,31.24609375,0.0,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,base,2,False,3.82,91.96875,0.670958,31.24609375,31.24609375,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,base,3,False,4.315,93.21875,0.844393,31.24609375,31.24609375,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,pr,3,False,3.581,92.390625,0.518655,31.24609375,0.0,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,pr,4,False,3.517,91.140625,0.583678,31.24609375,0.0,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,base,4,False,3.728,91.46875,0.592963,31.24609375,31.24609375,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,base,5,False,3.787,90.671875,0.746462,31.24609375,31.24609375,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_1task,pr,5,False,3.44,91.015625,0.55728,31.24609375,0.0,0.3662109375,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,base,warmup,True,50.328,462.375,13.0400635,249.96875,249.96875,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,pr,warmup,True,52.514,442.296875,12.6570325,249.96875,0.0,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,base,1,False,56.132,412.78125,8.1693925,249.96875,249.96875,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,pr,1,False,62.168,398.53125,10.9042115,249.96875,0.0,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,pr,2,False,55.738,393.203125,9.688614999999999,249.96875,0.0,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,base,2,False,55.412,401.703125,14.013288500000002,249.96875,249.96875,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,base,3,False,52.049,418.984375,8.771998,249.96875,249.96875,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   
16000p_8tasks,pr,3,False,51.212,475.609375,14.9291575,249.96875,0.0,2.9296875,88,2883584,412046047,8cc3776a5a0b6100f43389b98397541fe3b22e2713e46f1c43949dcb756eda60
   ```
   
   Benchmark-only instrumentation:
   
   ```diff
   diff --git a/native/shuffle/src/bin/shuffle_bench.rs 
b/native/shuffle/src/bin/shuffle_bench.rs
   index c2c9db30f..fa6fd4d45 100644
   --- a/native/shuffle/src/bin/shuffle_bench.rs
   +++ b/native/shuffle/src/bin/shuffle_bench.rs
   @@ -60,6 +60,9 @@ use std::time::Instant;
        about = "Standalone benchmark for Comet shuffle write performance"
    )]
    struct Args {
   +    #[arg(long)]
   +    validate_output: Option<PathBuf>,
   +
        /// Path to input Parquet file or directory of Parquet files
        #[arg(long)]
        input: PathBuf,
   @@ -124,6 +127,12 @@ struct Args {
    
    fn main() {
        let args = Args::parse();
   +    if let Some(ref output) = args.validate_output {
   +        let (_, rows) = read_parquet_metadata(&args.input, args.limit);
   +        validate_output(output, rows as usize);
   +        return;
   +    }
   +
    
        // Create output directory
        fs::create_dir_all(&args.output_dir).expect("Failed to create output 
directory");
   @@ -252,7 +261,9 @@ fn main() {
            }
        }
    
   -    let _ = fs::remove_file(&data_file);
   +    if std::env::var_os("COMET_BENCH_RETAIN_OUTPUT").is_none() {
   +        let _ = fs::remove_file(&data_file);
   +    }
    }
    
    fn print_shuffle_metrics(metrics: &MetricsSet, total_wall_time_secs: f64) {
   @@ -442,7 +453,7 @@ async fn execute_shuffle_write(
        limit: usize,
        data_file: String,
    ) -> datafusion::common::Result<(MetricsSet, MetricsSet)> {
   -    let config = SessionConfig::new().with_batch_size(batch_size);
   +    let config = 
SessionConfig::new().with_batch_size(batch_size).with_target_partitions(1);
        let mut runtime_builder = RuntimeEnvBuilder::new();
        if let Some(mem_limit) = memory_limit {
            runtime_builder = runtime_builder.with_memory_limit(mem_limit, 1.0);
   @@ -478,7 +489,7 @@ async fn execute_shuffle_write(
            input,
            partitioning,
            codec,
   -        data_file,
   +        data_file.clone(),
            false,
            write_buffer_size,
            max_buffer_bytes,
   @@ -489,6 +500,13 @@ async fn execute_shuffle_write(
        let stream = exec.execute(0, task_ctx).unwrap();
        collect(stream).await.unwrap();
    
   +    let metrics = exec.metrics().unwrap();
   +    let count = |name: &str| metrics.iter().filter(|m| m.value().name() == 
name).map(|m| m.value().as_usize()).sum::<usize>();
   +    println!("BENCH_SHUFFLE spill_count={} memory_spilled_bytes={} 
output_rows={}",
   +        count("spill_count"), count("memory_spilled_bytes"), 
count("output_rows"));
   +    let offsets = exec.partition_offsets().unwrap().get().unwrap();
   +    fs::write(Path::new(&data_file).with_extension("offsets.json"), 
format!("{offsets:?}")).unwrap();
   +
        // Collect metrics from the input plan (Parquet scan + optional 
coalesce)
        let input_metrics = collect_input_metrics(&exec);
    
   @@ -570,7 +588,9 @@ fn run_concurrent_shuffle_writes(
    
            for task_id in 0..args.concurrent_tasks {
                let task_dir = args.output_dir.join(format!("task_{task_id}"));
   -            let _ = fs::remove_dir_all(&task_dir);
   +            if std::env::var_os("COMET_BENCH_RETAIN_OUTPUT").is_none() {
   +                let _ = fs::remove_dir_all(&task_dir);
   +            }
            }
    
            start.elapsed().as_secs_f64()
   @@ -680,3 +700,48 @@ fn format_bytes(bytes: usize) -> String {
            format!("{bytes} B")
        }
    }
   +
   +// Fixture-specific readback runs in a separate process, outside the 
measured write.
   +fn validate_output(path: &Path, expected_rows: usize) {
   +    use arrow::array::Int64Array;
   +    use std::io::{BufReader, Read};
   +    let offsets: Vec<u64> = 
fs::read_to_string(path.with_extension("offsets.json")).unwrap()
   +        .trim().trim_matches(|c| c == '[' || c == ']')
   +        .split(',').map(|s| s.trim().parse().unwrap()).collect();
   +    assert_eq!(offsets[0], 0);
   +    assert_eq!(*offsets.last().unwrap(), fs::metadata(path).unwrap().len());
   +    let mut reader = BufReader::new(fs::File::open(path).unwrap());
   +    let mut position = 0u64;
   +    let mut seen = vec![0u8; expected_rows];
   +    let mut rows = 0usize;
   +    let mut blocks = 0usize;
   +    for pair in offsets.windows(2) {
   +        assert_eq!(position, pair[0]);
   +        assert!(pair[1] >= pair[0]);
   +        while position < pair[1] {
   +            let mut header = [0u8; 8];
   +            reader.read_exact(&mut header).unwrap();
   +            let len = u64::from_le_bytes(header) as usize;
   +            assert!(len >= 12 && len < 1024 * 1024 * 1024);
   +            let mut bytes = vec![0; len];
   +            reader.read_exact(&mut bytes).unwrap();
   +            let batch = 
datafusion_comet_shuffle::read_ipc_compressed(&bytes[8..]).unwrap();
   +            let ids = 
batch.column(0).as_any().downcast_ref::<Int64Array>().unwrap();
   +            let values = 
batch.column(1).as_any().downcast_ref::<Int64Array>().unwrap();
   +            for row in 0..batch.num_rows() {
   +                let id = ids.value(row);
   +                assert!(id >= 0 && (id as usize) < expected_rows);
   +                assert_eq!(seen[id as usize], 0, "duplicate row {id}");
   +                seen[id as usize] = 1;
   +                assert_eq!(values.value(row), id * 17 + 11);
   +                rows += 1;
   +            }
   +            position += 8 + len as u64;
   +            assert!(position <= pair[1], "block crosses a partition 
boundary");
   +            blocks += 1;
   +        }
   +    }
   +    assert_eq!(rows, expected_rows);
   +    assert!(seen.iter().all(|v| *v == 1));
   +    println!("BENCH_VALIDATED rows={rows} blocks={blocks} partitions={}", 
offsets.len() - 1);
   +}
   diff --git a/native/shuffle/src/writers/local/local_partition_writer.rs 
b/native/shuffle/src/writers/local/local_partition_writer.rs
   index d76ff23ef..6551f26fb 100644
   --- a/native/shuffle/src/writers/local/local_partition_writer.rs
   +++ b/native/shuffle/src/writers/local/local_partition_writer.rs
   @@ -94,6 +94,7 @@ pub(crate) struct LocalPartitionWriter {
        /// Id of the last partition passed to `finish_partition`, used to 
assert
        /// partitions are finalized in ascending order. `-1` before any call.
        last_finish_pid: i32,
   +    benchmark_merge_start: Option<std::time::Instant>,
    }
    
    impl LocalPartitionWriter {
   @@ -149,6 +150,7 @@ impl LocalPartitionWriter {
                write_buffer_size,
                num_output_partitions,
                last_finish_pid: -1,
   +            benchmark_merge_start: None,
            })
        }
    
   @@ -235,6 +237,13 @@ impl PartitionWriter for LocalPartitionWriter {
        where
            I: Iterator<Item = datafusion::common::Result<RecordBatch>>,
        {
   +        if pid == 0 {
   +            if let DataOutput::Multi { spill, .. } = &self.data_output {
   +                let (inner, outer) = spill.benchmark_range_capacity();
   +                println!("BENCH_METADATA before_merge inner_bytes={inner} 
outer_bytes={outer}");
   +            }
   +            self.benchmark_merge_start = Some(std::time::Instant::now());
   +        }
            if pid as i32 - self.last_finish_pid != 1 {
                return Err(DataFusionError::Execution(
                    "LocalPartitionWriter::finish_partition must be called in 
order.".to_string(),
   @@ -326,6 +335,13 @@ impl PartitionWriter for LocalPartitionWriter {
                    result.inspect_err(|_| recycled_buffer.clear())?;
                }
            }
   +        if pid + 1 == self.num_output_partitions {
   +            if let DataOutput::Multi { spill, .. } = &self.data_output {
   +                let elapsed = 
self.benchmark_merge_start.take().unwrap().elapsed().as_secs_f64();
   +                let (inner, outer) = spill.benchmark_range_capacity();
   +                println!("BENCH_METADATA after_merge inner_bytes={inner} 
outer_bytes={outer} merge_secs={elapsed:.6}");
   +            }
   +        }
            Ok(())
        }
    
   diff --git a/native/shuffle/src/writers/local/spill.rs 
b/native/shuffle/src/writers/local/spill.rs
   index ae4b8adfa..d6d4d235e 100644
   --- a/native/shuffle/src/writers/local/spill.rs
   +++ b/native/shuffle/src/writers/local/spill.rs
   @@ -69,6 +69,11 @@ pub(crate) struct PartitionedSpill {
    }
    
    impl PartitionedSpill {
   +    pub(crate) fn benchmark_range_capacity(&self) -> (usize, usize) {
   +        (self.ranges.iter().map(|v| v.capacity() * 
size_of::<Range<u64>>()).sum(),
   +         self.ranges.capacity() * size_of::<Vec<Range<u64>>>())
   +    }
   +
        pub(crate) fn new(
            shuffle_block_writer: ShuffleBlockWriter,
            write_buffer_size: usize,
   ```
   
   </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