From 847a7ce8f84ab058044f9ad1d82fc7095966aae2 Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Mon, 20 Jul 2026 13:20:13 +0800 Subject: [PATCH 1/3] perf(ipc): benchmark stream encoder StreamEncoder benchmarks catch regressions without measuring sink IO. Closes #10373 Signed-off-by: Jiawei Zhao --- arrow-ipc/benches/ipc_writer.rs | 72 ++++++++++++++++++++++++++++++++- 1 file changed, 71 insertions(+), 1 deletion(-) diff --git a/arrow-ipc/benches/ipc_writer.rs b/arrow-ipc/benches/ipc_writer.rs index 5050ff6cd112..bb4d819a8847 100644 --- a/arrow-ipc/benches/ipc_writer.rs +++ b/arrow-ipc/benches/ipc_writer.rs @@ -20,10 +20,14 @@ use arrow_array::builder::{ Date32Builder, Decimal128Builder, Int32Builder, StringBuilder, StringDictionaryBuilder, }; use arrow_array::types::UInt32Type; +use arrow_buffer::Buffer; use arrow_ipc::CompressionType; -use arrow_ipc::writer::{DictionaryHandling, FileWriter, IpcWriteOptions, StreamWriter}; +use arrow_ipc::writer::{ + DictionaryHandling, FileWriter, IpcWriteOptions, StreamEncoder, StreamWriter, +}; use arrow_schema::{DataType, Field, Schema}; use criterion::{Criterion, criterion_group, criterion_main}; +use std::hint::black_box; use std::sync::Arc; fn criterion_benchmark(c: &mut Criterion) { @@ -60,6 +64,36 @@ fn criterion_benchmark(c: &mut Criterion) { }) }); + group.bench_function("StreamEncoder/encode_10", |b| { + let batch = create_batch(8192, true); + b.iter(move || { + let mut encoder = StreamEncoder::try_new(batch.schema().as_ref()).unwrap(); + let mut encoded_len = 0; + for _ in 0..10 { + encoded_len += encoded_buffers_len(encoder.encode(&batch).unwrap()); + } + encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + black_box(encoded_len); + }) + }); + + group.bench_function("StreamEncoder/encode_10/zstd", |b| { + let batch = create_batch(8192, true); + b.iter(move || { + let options = IpcWriteOptions::default() + .try_with_compression(Some(CompressionType::ZSTD)) + .unwrap(); + let mut encoder = + StreamEncoder::try_new_with_options(batch.schema().as_ref(), options).unwrap(); + let mut encoded_len = 0; + for _ in 0..10 { + encoded_len += encoded_buffers_len(encoder.encode(&batch).unwrap()); + } + encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + black_box(encoded_len); + }) + }); + group.bench_function("FileWriter/write_10", |b| { let batch = create_batch(8192, true); let mut buffer = Vec::with_capacity(2 * 1024 * 1024); @@ -87,6 +121,20 @@ fn criterion_benchmark(c: &mut Criterion) { }) }); + group.bench_function("StreamEncoder/encode_10/dict", |b| { + let batches = create_unique_dict_batches(10, 8192); + let schema = batches[0].schema(); + b.iter(move || { + let mut encoder = StreamEncoder::try_new(schema.as_ref()).unwrap(); + let mut encoded_len = 0; + for batch in &batches { + encoded_len += encoded_buffers_len(encoder.encode(batch).unwrap()); + } + encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + black_box(encoded_len); + }) + }); + group.bench_function("StreamWriter/write_10/dict/delta", |b| { let batches = create_delta_dict_batches(10, 8192); let schema = batches[0].schema(); @@ -109,6 +157,24 @@ fn criterion_benchmark(c: &mut Criterion) { }) }); + group.bench_function("StreamEncoder/encode_10/dict/delta", |b| { + let batches = create_delta_dict_batches(10, 8192); + let schema = batches[0].schema(); + let options = + IpcWriteOptions::default().with_dictionary_handling(DictionaryHandling::Delta); + + b.iter(move || { + let mut encoder = + StreamEncoder::try_new_with_options(schema.as_ref(), options.clone()).unwrap(); + let mut encoded_len = 0; + for batch in &batches { + encoded_len += encoded_buffers_len(encoder.encode(batch).unwrap()); + } + encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + black_box(encoded_len); + }) + }); + // The file writer rejects dictionary replacement, so only the delta case is // exercised here (growing dictionaries that are prefixes of one another). group.bench_function("FileWriter/write_10/dict/delta", |b| { @@ -134,6 +200,10 @@ fn criterion_benchmark(c: &mut Criterion) { }); } +fn encoded_buffers_len(buffers: Vec) -> usize { + buffers.iter().map(Buffer::len).sum() +} + /// Build `n` record batches with a single dictionary column whose dictionary /// grows across batches. A single builder is reused with `finish_preserve_values` /// so each batch's dictionary has the previous batch's as a prefix which allows From 868d6439c796f7c87a7d2b95185432906db440fe Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Wed, 29 Jul 2026 11:04:36 +0800 Subject: [PATCH 2/3] perf(ipc): black-box encoded buffers Black-box StreamEncoder benchmark outputs before measuring their lengths so the benchmark boundary is less fragile. Signed-off-by: Jiawei Zhao --- arrow-ipc/benches/ipc_writer.rs | 24 ++++++++++++++++-------- 1 file changed, 16 insertions(+), 8 deletions(-) diff --git a/arrow-ipc/benches/ipc_writer.rs b/arrow-ipc/benches/ipc_writer.rs index bb4d819a8847..2db41014c11e 100644 --- a/arrow-ipc/benches/ipc_writer.rs +++ b/arrow-ipc/benches/ipc_writer.rs @@ -70,9 +70,11 @@ fn criterion_benchmark(c: &mut Criterion) { let mut encoder = StreamEncoder::try_new(batch.schema().as_ref()).unwrap(); let mut encoded_len = 0; for _ in 0..10 { - encoded_len += encoded_buffers_len(encoder.encode(&batch).unwrap()); + let encoded = black_box(encoder.encode(&batch).unwrap()); + encoded_len += encoded_buffers_len(encoded); } - encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + let encoded = black_box(encoder.finish().unwrap()); + encoded_len += encoded_buffers_len(encoded); black_box(encoded_len); }) }); @@ -87,9 +89,11 @@ fn criterion_benchmark(c: &mut Criterion) { StreamEncoder::try_new_with_options(batch.schema().as_ref(), options).unwrap(); let mut encoded_len = 0; for _ in 0..10 { - encoded_len += encoded_buffers_len(encoder.encode(&batch).unwrap()); + let encoded = black_box(encoder.encode(&batch).unwrap()); + encoded_len += encoded_buffers_len(encoded); } - encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + let encoded = black_box(encoder.finish().unwrap()); + encoded_len += encoded_buffers_len(encoded); black_box(encoded_len); }) }); @@ -128,9 +132,11 @@ fn criterion_benchmark(c: &mut Criterion) { let mut encoder = StreamEncoder::try_new(schema.as_ref()).unwrap(); let mut encoded_len = 0; for batch in &batches { - encoded_len += encoded_buffers_len(encoder.encode(batch).unwrap()); + let encoded = black_box(encoder.encode(batch).unwrap()); + encoded_len += encoded_buffers_len(encoded); } - encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + let encoded = black_box(encoder.finish().unwrap()); + encoded_len += encoded_buffers_len(encoded); black_box(encoded_len); }) }); @@ -168,9 +174,11 @@ fn criterion_benchmark(c: &mut Criterion) { StreamEncoder::try_new_with_options(schema.as_ref(), options.clone()).unwrap(); let mut encoded_len = 0; for batch in &batches { - encoded_len += encoded_buffers_len(encoder.encode(batch).unwrap()); + let encoded = black_box(encoder.encode(batch).unwrap()); + encoded_len += encoded_buffers_len(encoded); } - encoded_len += encoded_buffers_len(encoder.finish().unwrap()); + let encoded = black_box(encoder.finish().unwrap()); + encoded_len += encoded_buffers_len(encoded); black_box(encoded_len); }) }); From 835f3a243e6436ffbad7847a6724b6aafa0d765d Mon Sep 17 00:00:00 2001 From: Jiawei Zhao Date: Thu, 30 Jul 2026 20:48:38 +0800 Subject: [PATCH 3/3] perf(ipc): simplify encoder bench Black-boxing the StreamEncoder outputs is enough to keep the work observable. Avoid traversing buffer lengths in the measured path. Signed-off-by: Jiawei Zhao --- arrow-ipc/benches/ipc_writer.rs | 37 +++++++-------------------------- 1 file changed, 8 insertions(+), 29 deletions(-) diff --git a/arrow-ipc/benches/ipc_writer.rs b/arrow-ipc/benches/ipc_writer.rs index 2db41014c11e..5364049fb9f0 100644 --- a/arrow-ipc/benches/ipc_writer.rs +++ b/arrow-ipc/benches/ipc_writer.rs @@ -20,7 +20,6 @@ use arrow_array::builder::{ Date32Builder, Decimal128Builder, Int32Builder, StringBuilder, StringDictionaryBuilder, }; use arrow_array::types::UInt32Type; -use arrow_buffer::Buffer; use arrow_ipc::CompressionType; use arrow_ipc::writer::{ DictionaryHandling, FileWriter, IpcWriteOptions, StreamEncoder, StreamWriter, @@ -68,14 +67,10 @@ fn criterion_benchmark(c: &mut Criterion) { let batch = create_batch(8192, true); b.iter(move || { let mut encoder = StreamEncoder::try_new(batch.schema().as_ref()).unwrap(); - let mut encoded_len = 0; for _ in 0..10 { - let encoded = black_box(encoder.encode(&batch).unwrap()); - encoded_len += encoded_buffers_len(encoded); + black_box(encoder.encode(&batch).unwrap()); } - let encoded = black_box(encoder.finish().unwrap()); - encoded_len += encoded_buffers_len(encoded); - black_box(encoded_len); + black_box(encoder.finish().unwrap()); }) }); @@ -87,14 +82,10 @@ fn criterion_benchmark(c: &mut Criterion) { .unwrap(); let mut encoder = StreamEncoder::try_new_with_options(batch.schema().as_ref(), options).unwrap(); - let mut encoded_len = 0; for _ in 0..10 { - let encoded = black_box(encoder.encode(&batch).unwrap()); - encoded_len += encoded_buffers_len(encoded); + black_box(encoder.encode(&batch).unwrap()); } - let encoded = black_box(encoder.finish().unwrap()); - encoded_len += encoded_buffers_len(encoded); - black_box(encoded_len); + black_box(encoder.finish().unwrap()); }) }); @@ -130,14 +121,10 @@ fn criterion_benchmark(c: &mut Criterion) { let schema = batches[0].schema(); b.iter(move || { let mut encoder = StreamEncoder::try_new(schema.as_ref()).unwrap(); - let mut encoded_len = 0; for batch in &batches { - let encoded = black_box(encoder.encode(batch).unwrap()); - encoded_len += encoded_buffers_len(encoded); + black_box(encoder.encode(batch).unwrap()); } - let encoded = black_box(encoder.finish().unwrap()); - encoded_len += encoded_buffers_len(encoded); - black_box(encoded_len); + black_box(encoder.finish().unwrap()); }) }); @@ -172,14 +159,10 @@ fn criterion_benchmark(c: &mut Criterion) { b.iter(move || { let mut encoder = StreamEncoder::try_new_with_options(schema.as_ref(), options.clone()).unwrap(); - let mut encoded_len = 0; for batch in &batches { - let encoded = black_box(encoder.encode(batch).unwrap()); - encoded_len += encoded_buffers_len(encoded); + black_box(encoder.encode(batch).unwrap()); } - let encoded = black_box(encoder.finish().unwrap()); - encoded_len += encoded_buffers_len(encoded); - black_box(encoded_len); + black_box(encoder.finish().unwrap()); }) }); @@ -208,10 +191,6 @@ fn criterion_benchmark(c: &mut Criterion) { }); } -fn encoded_buffers_len(buffers: Vec) -> usize { - buffers.iter().map(Buffer::len).sum() -} - /// Build `n` record batches with a single dictionary column whose dictionary /// grows across batches. A single builder is reused with `finish_preserve_values` /// so each batch's dictionary has the previous batch's as a prefix which allows