diff --git a/arrow-ipc/benches/ipc_writer.rs b/arrow-ipc/benches/ipc_writer.rs index 5050ff6cd112..5364049fb9f0 100644 --- a/arrow-ipc/benches/ipc_writer.rs +++ b/arrow-ipc/benches/ipc_writer.rs @@ -21,9 +21,12 @@ use arrow_array::builder::{ }; use arrow_array::types::UInt32Type; 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 +63,32 @@ 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(); + for _ in 0..10 { + black_box(encoder.encode(&batch).unwrap()); + } + black_box(encoder.finish().unwrap()); + }) + }); + + 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(); + for _ in 0..10 { + black_box(encoder.encode(&batch).unwrap()); + } + black_box(encoder.finish().unwrap()); + }) + }); + group.bench_function("FileWriter/write_10", |b| { let batch = create_batch(8192, true); let mut buffer = Vec::with_capacity(2 * 1024 * 1024); @@ -87,6 +116,18 @@ 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(); + for batch in &batches { + black_box(encoder.encode(batch).unwrap()); + } + black_box(encoder.finish().unwrap()); + }) + }); + group.bench_function("StreamWriter/write_10/dict/delta", |b| { let batches = create_delta_dict_batches(10, 8192); let schema = batches[0].schema(); @@ -109,6 +150,22 @@ 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(); + for batch in &batches { + black_box(encoder.encode(batch).unwrap()); + } + black_box(encoder.finish().unwrap()); + }) + }); + // 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| {