Skip to content

Commit 7abb225

Browse files
HippoBaroalamb
andauthored
bench(parquet): add ListArray benchmarks for runtime and peak memory (#9846)
# Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. --> - Contributes to #9731 - Dependency of #9848 # Rationale for this change See #9848 Existing benchmarks have some gaps in the types of columns they exercise. Additionally, I would like to improve the memory efficiency of the read/decode path in terms of RSS requirements, especially for sparse inputs and we currently do not have any infrastructure to measure that. # What changes are included in this PR? Extend the existing `arrow_reader` runtime benchmarks with `Int32` and `FixedBinary32` list columns alongside the existing `StringList`, with parameterized null density (0%, 50%, 90%, 99%). The prior benchmarks only covered string lists, which didn't surface costs specific to fixed-width and primitive element types. Add a new `arrow_reader_peak_memory` benchmark that measures peak heap usage during `ListArrayReader::consume_batch` using a thread-local tracking allocator. It captures how RSS-efficient we are when materializing a column into its final Arrow in-memory representation. # Are these changes tested? All tests passing. # Are there any user-facing changes? None. Signed-off-by: Hippolyte Barraud <hippolyte.barraud@datadoghq.com> Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
1 parent 6ce4bc8 commit 7abb225

3 files changed

Lines changed: 808 additions & 15 deletions

File tree

parquet/Cargo.toml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -289,5 +289,10 @@ required-features = ["arrow"]
289289
name = "bloom_filter"
290290
harness = false
291291

292+
[[bench]]
293+
name = "arrow_reader_peak_memory"
294+
required-features = ["arrow", "test_common", "experimental"]
295+
harness = false
296+
292297
[lib]
293298
bench = false

parquet/benches/arrow_reader.rs

Lines changed: 206 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,16 @@ fn build_test_schema() -> SchemaDescPtr {
9595
OPTIONAL GROUP optional_struct_optional_int32_leaf {
9696
OPTIONAL INT32 element;
9797
}
98+
OPTIONAL GROUP int32_list (LIST) {
99+
repeated group list {
100+
optional INT32 element;
101+
}
102+
}
103+
OPTIONAL GROUP fixed32_list (LIST) {
104+
repeated group list {
105+
optional FIXED_LEN_BYTE_ARRAY(32) element;
106+
}
107+
}
98108
}
99109
";
100110
parse_message_type(message_type)
@@ -668,6 +678,150 @@ fn build_string_list_page_iterator(
668678
InMemoryPageIterator::new(pages)
669679
}
670680

681+
fn build_int32_list_page_iterator(
682+
column_desc: ColumnDescPtr,
683+
null_density: f32,
684+
) -> impl PageIterator + Clone {
685+
let max_def_level = column_desc.max_def_level();
686+
let max_rep_level = column_desc.max_rep_level();
687+
assert_eq!(max_def_level, 3);
688+
assert_eq!(max_rep_level, 1);
689+
690+
let mut rng = seedable_rng();
691+
let mut pages: Vec<Vec<parquet::column::page::Page>> = Vec::new();
692+
for _i in 0..NUM_ROW_GROUPS {
693+
let mut column_chunk_pages = Vec::new();
694+
for _j in 0..PAGES_PER_GROUP {
695+
let mut values: Vec<i32> = Vec::with_capacity(VALUES_PER_PAGE * MAX_LIST_LEN);
696+
let mut def_levels = Vec::with_capacity(VALUES_PER_PAGE * MAX_LIST_LEN);
697+
let mut rep_levels = Vec::with_capacity(VALUES_PER_PAGE * MAX_LIST_LEN);
698+
for _k in 0..VALUES_PER_PAGE {
699+
rep_levels.push(0);
700+
if rng.random::<f32>() < null_density {
701+
def_levels.push(0);
702+
continue;
703+
}
704+
let len = rng.random_range(0..MAX_LIST_LEN);
705+
if len == 0 {
706+
def_levels.push(1);
707+
continue;
708+
}
709+
710+
(1..len).for_each(|_| rep_levels.push(1));
711+
712+
for _l in 0..len {
713+
if rng.random::<f32>() < null_density {
714+
def_levels.push(2);
715+
} else {
716+
def_levels.push(3);
717+
values.push(rng.random());
718+
}
719+
}
720+
}
721+
let mut page_builder =
722+
DataPageBuilderImpl::new(column_desc.clone(), values.len() as u32, true);
723+
page_builder.add_rep_levels(max_rep_level, &rep_levels);
724+
page_builder.add_def_levels(max_def_level, &def_levels);
725+
page_builder.add_values::<Int32Type>(Encoding::PLAIN, &values);
726+
column_chunk_pages.push(page_builder.consume());
727+
}
728+
pages.push(column_chunk_pages);
729+
}
730+
731+
InMemoryPageIterator::new(pages)
732+
}
733+
734+
fn create_int32_list_reader(
735+
page_iterator: impl PageIterator + 'static,
736+
column_desc: ColumnDescPtr,
737+
) -> Box<dyn ArrayReader> {
738+
use parquet::arrow::array_reader::PrimitiveArrayReader;
739+
let items = Box::new(
740+
PrimitiveArrayReader::<Int32Type>::new(
741+
Box::new(page_iterator),
742+
column_desc,
743+
None,
744+
DEFAULT_BATCH_SIZE,
745+
)
746+
.unwrap(),
747+
) as Box<dyn ArrayReader>;
748+
let field = Field::new_list_field(DataType::Int32, true);
749+
let data_type = DataType::List(Arc::new(field));
750+
Box::new(ListArrayReader::<i32>::new(items, data_type, 2, 1, true))
751+
}
752+
753+
const FIXED_BYTE_LEN: usize = 32;
754+
755+
fn build_fixed32_list_page_iterator(
756+
column_desc: ColumnDescPtr,
757+
null_density: f32,
758+
) -> impl PageIterator + Clone {
759+
let max_def_level = column_desc.max_def_level();
760+
let max_rep_level = column_desc.max_rep_level();
761+
assert_eq!(max_def_level, 3);
762+
assert_eq!(max_rep_level, 1);
763+
764+
let mut rng = seedable_rng();
765+
let mut pages: Vec<Vec<parquet::column::page::Page>> = Vec::new();
766+
for _i in 0..NUM_ROW_GROUPS {
767+
let mut column_chunk_pages = Vec::new();
768+
for _j in 0..PAGES_PER_GROUP {
769+
let mut values: Vec<parquet::data_type::FixedLenByteArray> =
770+
Vec::with_capacity(VALUES_PER_PAGE * MAX_LIST_LEN);
771+
let mut def_levels = Vec::with_capacity(VALUES_PER_PAGE * MAX_LIST_LEN);
772+
let mut rep_levels = Vec::with_capacity(VALUES_PER_PAGE * MAX_LIST_LEN);
773+
for _k in 0..VALUES_PER_PAGE {
774+
rep_levels.push(0);
775+
if rng.random::<f32>() < null_density {
776+
def_levels.push(0);
777+
continue;
778+
}
779+
let len = rng.random_range(0..MAX_LIST_LEN);
780+
if len == 0 {
781+
def_levels.push(1);
782+
continue;
783+
}
784+
(1..len).for_each(|_| rep_levels.push(1));
785+
for _l in 0..len {
786+
if rng.random::<f32>() < null_density {
787+
def_levels.push(2);
788+
} else {
789+
def_levels.push(3);
790+
let mut buf = vec![0u8; FIXED_BYTE_LEN];
791+
rng.fill(&mut buf[..]);
792+
values.push(buf.into());
793+
}
794+
}
795+
}
796+
let mut page_builder =
797+
DataPageBuilderImpl::new(column_desc.clone(), values.len() as u32, true);
798+
page_builder.add_rep_levels(max_rep_level, &rep_levels);
799+
page_builder.add_def_levels(max_def_level, &def_levels);
800+
page_builder.add_values::<FixedLenByteArrayType>(Encoding::PLAIN, &values);
801+
column_chunk_pages.push(page_builder.consume());
802+
}
803+
pages.push(column_chunk_pages);
804+
}
805+
806+
InMemoryPageIterator::new(pages)
807+
}
808+
809+
fn create_fixed32_list_reader(
810+
page_iterator: impl PageIterator + 'static,
811+
column_desc: ColumnDescPtr,
812+
) -> Box<dyn ArrayReader> {
813+
let items = make_fixed_len_byte_array_reader(
814+
Box::new(page_iterator),
815+
column_desc,
816+
None,
817+
DEFAULT_BATCH_SIZE,
818+
)
819+
.unwrap();
820+
let field = Field::new_list_field(DataType::FixedSizeBinary(FIXED_BYTE_LEN as i32), true);
821+
let data_type = DataType::List(Arc::new(field));
822+
Box::new(ListArrayReader::<i32>::new(items, data_type, 2, 1, true))
823+
}
824+
671825
fn bench_array_reader(mut array_reader: Box<dyn ArrayReader>) -> usize {
672826
// test procedure: read data in batches of 8192 until no more data
673827
let mut total_count = 0;
@@ -1645,6 +1799,8 @@ fn add_benches(c: &mut Criterion) {
16451799
let optional_uint64_column_desc = schema.column(38);
16461800
let mandatory_struct_optional_in32_column_desc = schema.column(39);
16471801
let optional_struct_optional_in32_column_desc = schema.column(40);
1802+
let int32_list_desc = schema.column(41);
1803+
let fixed32_list_desc = schema.column(42);
16481804

16491805
// primitive / int32 benchmarks
16501806
// =============================
@@ -2228,24 +2384,59 @@ fn add_benches(c: &mut Criterion) {
22282384
// list benchmarks
22292385
//==============================
22302386

2231-
let list_data = build_string_list_page_iterator(string_list_desc.clone(), 0.);
2232-
let mut group = c.benchmark_group("arrow_array_reader/ListArray");
2233-
group.bench_function("plain encoded optional strings no NULLs", |b| {
2234-
b.iter(|| {
2235-
let reader = create_string_list_reader(list_data.clone(), string_list_desc.clone());
2236-
count = bench_array_reader(reader);
2387+
let mut group = c.benchmark_group("arrow_array_reader/ListArray/StringList");
2388+
for (label, null_density) in [
2389+
("no NULLs", 0.0),
2390+
("half NULLs", 0.5),
2391+
("90pct NULLs", 0.9),
2392+
("99pct NULLs", 0.99),
2393+
] {
2394+
let list_data = build_string_list_page_iterator(string_list_desc.clone(), null_density);
2395+
group.bench_function(label, |b| {
2396+
b.iter(|| {
2397+
let reader = create_string_list_reader(list_data.clone(), string_list_desc.clone());
2398+
count = bench_array_reader(reader);
2399+
});
2400+
assert_eq!(count, EXPECTED_VALUE_COUNT);
22372401
});
2238-
assert_eq!(count, EXPECTED_VALUE_COUNT);
2239-
});
2240-
let list_data = build_string_list_page_iterator(string_list_desc.clone(), 0.5);
2241-
group.bench_function("plain encoded optional strings half NULLs", |b| {
2242-
b.iter(|| {
2243-
let reader = create_string_list_reader(list_data.clone(), string_list_desc.clone());
2244-
count = bench_array_reader(reader);
2402+
}
2403+
group.finish();
2404+
2405+
let mut group = c.benchmark_group("arrow_array_reader/ListArray/Int32List");
2406+
for (label, null_density) in [
2407+
("no NULLs", 0.0),
2408+
("half NULLs", 0.5),
2409+
("90pct NULLs", 0.9),
2410+
("99pct NULLs", 0.99),
2411+
] {
2412+
let list_data = build_int32_list_page_iterator(int32_list_desc.clone(), null_density);
2413+
group.bench_function(label, |b| {
2414+
b.iter(|| {
2415+
let reader = create_int32_list_reader(list_data.clone(), int32_list_desc.clone());
2416+
count = bench_array_reader(reader);
2417+
});
2418+
assert_eq!(count, EXPECTED_VALUE_COUNT);
22452419
});
2246-
assert_eq!(count, EXPECTED_VALUE_COUNT);
2247-
});
2420+
}
2421+
group.finish();
22482422

2423+
let mut group = c.benchmark_group("arrow_array_reader/ListArray/Fixed32List");
2424+
for (label, null_density) in [
2425+
("no NULLs", 0.0),
2426+
("half NULLs", 0.5),
2427+
("90pct NULLs", 0.9),
2428+
("99pct NULLs", 0.99),
2429+
] {
2430+
let list_data = build_fixed32_list_page_iterator(fixed32_list_desc.clone(), null_density);
2431+
group.bench_function(label, |b| {
2432+
b.iter(|| {
2433+
let reader =
2434+
create_fixed32_list_reader(list_data.clone(), fixed32_list_desc.clone());
2435+
count = bench_array_reader(reader);
2436+
});
2437+
assert_eq!(count, EXPECTED_VALUE_COUNT);
2438+
});
2439+
}
22492440
group.finish();
22502441

22512442
// fixed_len_byte_array benchmarks

0 commit comments

Comments
 (0)