diff --git a/vortex-arrow/benches/to_arrow.rs b/vortex-arrow/benches/to_arrow.rs index e6016ec9705..e5fca5b51f0 100644 --- a/vortex-arrow/benches/to_arrow.rs +++ b/vortex-arrow/benches/to_arrow.rs @@ -3,6 +3,8 @@ #![expect(clippy::unwrap_used)] +use std::fmt::Display; +use std::fmt::Formatter; use std::sync::Arc; use std::sync::LazyLock; @@ -10,6 +12,7 @@ use arrow_schema::DataType; use arrow_schema::Field; use divan::Bencher; use divan::counter::ItemsCount; +use itertools::iproduct; use vortex_array::ArrayRef; use vortex_array::ExecutionCtx; use vortex_array::IntoArray; @@ -18,7 +21,6 @@ use vortex_array::array_session; use vortex_array::arrays::ChunkedArray; use vortex_array::arrays::DecimalArray; use vortex_array::arrays::DictArray; -use vortex_array::arrays::FilterArray; use vortex_array::arrays::ListArray; use vortex_array::arrays::PrimitiveArray; use vortex_array::arrays::StructArray; @@ -41,7 +43,6 @@ use vortex_arrow::ArrowSessionExt; use vortex_arrow::dtype::ToArrowType as _; use vortex_fsst::fsst_compress; use vortex_fsst::fsst_train_compressor; -use vortex_mask::Mask; use vortex_onpair::DEFAULT_CONFIG; use vortex_onpair::onpair_compress; use vortex_session::VortexSession; @@ -121,205 +122,274 @@ fn ArrowExportVTable_to_arrow_field(bencher: Bencher) { .bench_values(|dtype| SESSION.arrow().to_arrow_field("", &dtype).unwrap()) } -#[derive(Clone, Copy, Debug)] +#[derive(Clone, Copy)] enum StringEncoding { + Offset, View, Fsst, OnPair, Zstd, +} + +impl Display for StringEncoding { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::Offset => "offset", + Self::View => "view", + Self::Fsst => "fsst", + Self::OnPair => "onpair", + Self::Zstd => "zstd", + }) + } +} + +#[derive(Clone, Copy)] +enum StringStructure { + Flat, Dict, - DictFsst, - DictZstd, - FilterFsst, - FilterZstd, - FilterDictFsst, - ChunkedFsst, - /// Every third row null, so the export walks a partial validity mask rather than an all-valid - /// one and has to interleave nulls with the decoded values. - NullableFsst, - NullableZstd, - NullableDict, + Chunked, +} + +impl Display for StringStructure { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::Flat => "flat", + Self::Dict => "dict", + Self::Chunked => "chunked", + }) + } +} + +#[derive(Clone, Copy)] +enum StringValidity { + NonNullable, + Nullable, +} + +impl StringValidity { + fn nullability(self) -> Nullability { + match self { + Self::NonNullable => Nullability::NonNullable, + Self::Nullable => Nullability::Nullable, + } + } +} + +impl Display for StringValidity { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::NonNullable => "nonnull", + Self::Nullable => "nullable", + }) + } +} + +#[derive(Clone, Copy)] +struct StringCase { + encoding: StringEncoding, + structure: StringStructure, + validity: StringValidity, +} + +impl Display for StringCase { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "{}/{}/{}", self.encoding, self.structure, self.validity) + } +} + +#[derive(Clone, Copy)] +enum ArrowStringLayout { + Offset, + View, +} + +impl ArrowStringLayout { + fn data_type(self) -> DataType { + match self { + Self::Offset => DataType::Utf8, + Self::View => DataType::Utf8View, + } + } +} + +impl Display for ArrowStringLayout { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::Offset => "offset", + Self::View => "view", + }) + } +} + +#[derive(Clone, Copy)] +struct StringExportCase { + array: StringCase, + layout: ArrowStringLayout, +} + +impl Display for StringExportCase { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "{}/{}", self.layout, self.array) + } } const STRING_ENCODINGS: &[StringEncoding] = &[ + StringEncoding::Offset, StringEncoding::View, StringEncoding::Fsst, StringEncoding::OnPair, StringEncoding::Zstd, - StringEncoding::Dict, - StringEncoding::DictFsst, - StringEncoding::DictZstd, - StringEncoding::FilterFsst, - StringEncoding::FilterZstd, - StringEncoding::FilterDictFsst, - StringEncoding::ChunkedFsst, - StringEncoding::NullableFsst, - StringEncoding::NullableZstd, - StringEncoding::NullableDict, ]; - -/// Encodings whose `append_to_builder` the builder benchmarks reach directly. -/// -/// The Arrow export cannot stand in for these: `execute_until` stops at the first canonical array, -/// so a bare FSST/OnPair/Zstd root is canonicalized to `VarBinView` before any builder sees it. -/// Only `Chunked`, `Constant` and `VarBin` roots reach an encoding's own `append_to_builder` that -/// way, whereas the scan machinery appends encoded arrays into a builder directly. -const BUILDER_STRING_ENCODINGS: &[StringEncoding] = &[ - StringEncoding::View, - StringEncoding::Fsst, - StringEncoding::OnPair, - StringEncoding::Zstd, - StringEncoding::Dict, - StringEncoding::ChunkedFsst, - StringEncoding::NullableFsst, - StringEncoding::NullableZstd, - StringEncoding::NullableDict, +const STRING_STRUCTURES: &[StringStructure] = &[ + StringStructure::Flat, + StringStructure::Dict, + StringStructure::Chunked, ]; +const STRING_VALIDITIES: &[StringValidity] = + &[StringValidity::NonNullable, StringValidity::Nullable]; +const ARROW_STRING_LAYOUTS: &[ArrowStringLayout] = + &[ArrowStringLayout::Offset, ArrowStringLayout::View]; -const OFFSET_STRING_ROWS: usize = 100_000; -const OFFSET_STRING_CHUNKS: usize = 4; +const STRING_ROWS: usize = 100_000; +const STRING_CHUNKS: usize = 4; const DICTIONARY_SIZE: usize = 2_048; -fn structured_strings(len: usize) -> VarBinViewArray { - let values = (0..len) - .map(|index| format!("https://example.com/common/path/{index:06}/shared-suffix")) - .collect::>(); - VarBinViewArray::from_iter_str(values.iter().map(String::as_str)) +fn string_cases() -> Vec { + iproduct!( + STRING_ENCODINGS.iter().copied(), + STRING_STRUCTURES.iter().copied(), + STRING_VALIDITIES.iter().copied() + ) + .map(|(encoding, structure, validity)| StringCase { + encoding, + structure, + validity, + }) + .collect() } -fn nullable_structured_strings(len: usize) -> VarBinViewArray { - let values = (0..len) - .map(|index| { - (!index.is_multiple_of(3)) - .then(|| format!("https://example.com/common/path/{index:06}/shared-suffix")) - }) - .collect::>(); - VarBinViewArray::from_iter( - values.iter().map(|value| value.as_deref()), - DType::Utf8(Nullability::Nullable), - ) +fn string_export_cases() -> Vec { + iproduct!(string_cases(), ARROW_STRING_LAYOUTS.iter().copied()) + .map(|(array, layout)| StringExportCase { array, layout }) + .collect() } -fn dictionary_values() -> VarBinViewArray { - structured_strings(DICTIONARY_SIZE) +fn structured_strings(len: usize, validity: StringValidity) -> VarBinViewArray { + match validity { + StringValidity::NonNullable => { + let values = (0..len) + .map(|index| format!("https://example.com/common/path/{index:06}/shared-suffix")) + .collect::>(); + VarBinViewArray::from_iter_str(values.iter().map(String::as_str)) + } + StringValidity::Nullable => { + let values = (0..len) + .map(|index| { + (!index.is_multiple_of(3)).then(|| { + format!("https://example.com/common/path/{index:06}/shared-suffix") + }) + }) + .collect::>(); + VarBinViewArray::from_iter( + values.iter().map(|value| value.as_deref()), + DType::Utf8(Nullability::Nullable), + ) + } + } } fn dictionary_codes() -> ArrayRef { PrimitiveArray::from_iter( - (0..OFFSET_STRING_ROWS).map(|index| u16::try_from(index % DICTIONARY_SIZE).unwrap()), + (0..STRING_ROWS).map(|index| u16::try_from(index % DICTIONARY_SIZE).unwrap()), ) .into_array() } -fn half_rows_mask() -> Mask { - Mask::from_iter((0..OFFSET_STRING_ROWS).map(|index| index.is_multiple_of(2))) +fn offset(source: VarBinViewArray, ctx: &mut ExecutionCtx) -> ArrayRef { + let source = source.into_array(); + let mut builder = VarBinBuilder::::with_capacity(source.dtype().clone(), source.len()); + source.append_to_builder(&mut builder, ctx).unwrap(); + builder.finish_into_varbin().into_array() } -fn filtered(array: ArrayRef) -> ArrayRef { - // Keep Filter as a lazy intermediate so benchmark setup cannot optimize it away. - FilterArray::new(array, half_rows_mask()).into_array() +fn fsst(source: ArrayRef, ctx: &mut ExecutionCtx) -> ArrayRef { + let compressor = fsst_train_compressor(&source, ctx).unwrap(); + fsst_compress(&source, &compressor, ctx) + .unwrap() + .into_array() } -fn chunked_fsst(ctx: &mut ExecutionCtx) -> ArrayRef { - let source = structured_strings(OFFSET_STRING_ROWS).into_array(); - let compressor = fsst_train_compressor(&source, ctx).unwrap(); - let chunk_size = OFFSET_STRING_ROWS / OFFSET_STRING_CHUNKS; - let chunks = (0..OFFSET_STRING_CHUNKS).map(|chunk_index| { - let start = chunk_index * chunk_size; - let end = if chunk_index + 1 == OFFSET_STRING_CHUNKS { - OFFSET_STRING_ROWS - } else { - start + chunk_size - }; - let chunk = source.slice(start..end).unwrap(); - fsst_compress(&chunk, &compressor, ctx) +fn encode_strings( + source: VarBinViewArray, + encoding: StringEncoding, + ctx: &mut ExecutionCtx, +) -> ArrayRef { + match encoding { + StringEncoding::Offset => offset(source, ctx), + StringEncoding::View => source.into_array(), + StringEncoding::Fsst => fsst(source.into_array(), ctx), + StringEncoding::OnPair => { + onpair_compress(&source.into_array(), DEFAULT_CONFIG, ctx).unwrap() + } + StringEncoding::Zstd => Zstd::from_var_bin_view_without_dict(&source, 3, 8_192, ctx) .unwrap() - .into_array() - }); - ChunkedArray::try_new(chunks, source.dtype().clone()) + .into_array(), + } +} + +fn dictionary_strings( + encoding: StringEncoding, + validity: StringValidity, + ctx: &mut ExecutionCtx, +) -> ArrayRef { + let values = encode_strings(structured_strings(DICTIONARY_SIZE, validity), encoding, ctx); + DictArray::try_new(dictionary_codes(), values) .unwrap() .into_array() } -fn fsst(source: ArrayRef, ctx: &mut ExecutionCtx) -> ArrayRef { - let compressor = fsst_train_compressor(&source, ctx).unwrap(); - fsst_compress(&source, &compressor, ctx) +fn chunked_strings( + encoding: StringEncoding, + validity: StringValidity, + ctx: &mut ExecutionCtx, +) -> ArrayRef { + let chunk_size = STRING_ROWS / STRING_CHUNKS; + let chunks = (0..STRING_CHUNKS).map(|chunk_index| { + let start = chunk_index * chunk_size; + let end = if chunk_index + 1 == STRING_CHUNKS { + STRING_ROWS + } else { + start + chunk_size + }; + encode_strings(structured_strings(end - start, validity), encoding, ctx) + }); + ChunkedArray::try_new(chunks, DType::Utf8(validity.nullability())) .unwrap() .into_array() } -fn string_array(encoding: StringEncoding) -> ArrayRef { +fn string_array(case: StringCase) -> ArrayRef { let mut ctx = SESSION.create_execution_ctx(); - match encoding { - StringEncoding::View => structured_strings(OFFSET_STRING_ROWS).into_array(), - StringEncoding::Fsst => fsst( - structured_strings(OFFSET_STRING_ROWS).into_array(), + match case.structure { + StringStructure::Flat => encode_strings( + structured_strings(STRING_ROWS, case.validity), + case.encoding, &mut ctx, ), - StringEncoding::OnPair => onpair_compress( - &structured_strings(OFFSET_STRING_ROWS).into_array(), - DEFAULT_CONFIG, - &mut ctx, - ) - .unwrap(), - StringEncoding::Zstd => { - let source = structured_strings(OFFSET_STRING_ROWS); - Zstd::from_var_bin_view_without_dict(&source, 3, 8_192, &mut ctx) - .unwrap() - .into_array() - } - StringEncoding::Dict => { - DictArray::try_new(dictionary_codes(), dictionary_values().into_array()) - .unwrap() - .into_array() - } - StringEncoding::DictFsst => { - let values = fsst(dictionary_values().into_array(), &mut ctx); - DictArray::try_new(dictionary_codes(), values) - .unwrap() - .into_array() - } - StringEncoding::DictZstd => { - let values = dictionary_values(); - let compressed_values = - Zstd::from_var_bin_view_without_dict(&values, 3, 8_192, &mut ctx) - .unwrap() - .into_array(); - DictArray::try_new(dictionary_codes(), compressed_values) - .unwrap() - .into_array() - } - StringEncoding::FilterFsst => filtered(string_array(StringEncoding::Fsst)), - StringEncoding::FilterZstd => filtered(string_array(StringEncoding::Zstd)), - StringEncoding::FilterDictFsst => filtered(string_array(StringEncoding::DictFsst)), - StringEncoding::ChunkedFsst => chunked_fsst(&mut ctx), - StringEncoding::NullableFsst => fsst( - nullable_structured_strings(OFFSET_STRING_ROWS).into_array(), - &mut ctx, - ), - StringEncoding::NullableZstd => { - let source = nullable_structured_strings(OFFSET_STRING_ROWS); - Zstd::from_var_bin_view_without_dict(&source, 3, 8_192, &mut ctx) - .unwrap() - .into_array() - } - StringEncoding::NullableDict => { - // Nulls live in the dictionary rather than the codes, so the export has to combine - // the two validities. - let values = nullable_structured_strings(DICTIONARY_SIZE); - DictArray::try_new(dictionary_codes(), values.into_array()) - .unwrap() - .into_array() - } + StringStructure::Dict => dictionary_strings(case.encoding, case.validity, &mut ctx), + StringStructure::Chunked => chunked_strings(case.encoding, case.validity, &mut ctx), } } -/// End-to-end export to Arrow `Utf8`, which is served through a `VarBinBuilder`. -#[divan::bench(args = STRING_ENCODINGS)] -fn offset_string_export(bencher: Bencher, encoding: StringEncoding) { - let array = string_array(encoding); - let field = Field::new("value", DataType::Utf8, array.dtype().is_nullable()); - +/// Measures export to Arrow offset arrays and Arrow view arrays. +#[divan::bench(args = string_export_cases())] +fn string_export(bencher: Bencher, case: StringExportCase) { + let array = string_array(case.array); + let field = Field::new( + "value", + case.layout.data_type(), + array.dtype().is_nullable(), + ); bencher .with_inputs(|| (array.clone(), SESSION.create_execution_ctx())) .input_counter(|(array, _)| ItemsCount::new(array.len())) @@ -331,10 +401,10 @@ fn offset_string_export(bencher: Bencher, encoding: StringEncoding) { }); } -/// Appends an encoded array straight into an offset builder. -#[divan::bench(args = BUILDER_STRING_ENCODINGS)] -fn append_to_varbin_builder(bencher: Bencher, encoding: StringEncoding) { - let array = string_array(encoding); +/// Measures a direct append to an offset builder. +#[divan::bench(args = string_cases())] +fn append_to_varbin_builder(bencher: Bencher, case: StringCase) { + let array = string_array(case); bencher .with_inputs(|| (array.clone(), SESSION.create_execution_ctx())) @@ -347,10 +417,10 @@ fn append_to_varbin_builder(bencher: Bencher, encoding: StringEncoding) { }); } -/// Appends an encoded array straight into a view builder. -#[divan::bench(args = BUILDER_STRING_ENCODINGS)] -fn append_to_view_builder(bencher: Bencher, encoding: StringEncoding) { - let array = string_array(encoding); +/// Measures a direct append to a view builder. +#[divan::bench(args = string_cases())] +fn append_to_view_builder(bencher: Bencher, case: StringCase) { + let array = string_array(case); bencher .with_inputs(|| (array.clone(), SESSION.create_execution_ctx()))