diff --git a/datafusion/functions-nested/benches/array_set_ops.rs b/datafusion/functions-nested/benches/array_set_ops.rs index d43bbdb577d0..fa5f26531aa6 100644 --- a/datafusion/functions-nested/benches/array_set_ops.rs +++ b/datafusion/functions-nested/benches/array_set_ops.rs @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -use arrow::array::{ArrayRef, Int64Array, ListArray}; +use arrow::array::{ArrayRef, Float64Array, Int64Array, ListArray}; use arrow::buffer::OffsetBuffer; use arrow::datatypes::{DataType, Field}; use criterion::{ @@ -38,6 +38,11 @@ const SEED: u64 = 42; /// Extra rows on each side when building sliced arrays, so the underlying /// values buffer is much larger than the visible portion. const SLICE_PADDING: usize = 5000; +/// Keep the visible slice at two values while varying the backing array from +/// 8 KiB to 8 MiB. +const SLICED_FLOAT_BACKING_VALUES: &[usize] = &[1024, 1024 * 1024]; +/// Keep the backing array at 8 MiB while varying the visible slice. +const SLICED_FLOAT_VISIBLE_VALUES: &[usize] = &[2, 2048]; fn criterion_benchmark(c: &mut Criterion) { bench_array_union(c); @@ -47,6 +52,7 @@ fn criterion_benchmark(c: &mut Criterion) { bench_array_union_sliced(c); bench_array_intersect_sliced(c); bench_array_distinct_sliced(c); + bench_array_distinct_sliced_float(c); bench_array_except_sliced(c); } @@ -69,6 +75,23 @@ fn invoke_udf(udf: &impl ScalarUDFImpl, array1: &ArrayRef, array2: &ArrayRef) { ); } +fn invoke_unary_udf( + udf: &impl ScalarUDFImpl, + array: &ArrayRef, + number_rows: usize, +) -> ColumnarValue { + black_box( + udf.invoke_with_args(ScalarFunctionArgs { + args: vec![ColumnarValue::Array(array.clone())], + arg_fields: vec![Field::new("arr", array.data_type().clone(), false).into()], + number_rows, + return_field: Field::new("result", array.data_type().clone(), false).into(), + config_options: Arc::new(ConfigOptions::default()), + }) + .unwrap(), + ) +} + fn bench_array_union(c: &mut Criterion) { let mut group = c.benchmark_group("array_union"); let udf = ArrayUnion::new(); @@ -139,28 +162,7 @@ fn bench_array_distinct(c: &mut Criterion) { group.bench_with_input( BenchmarkId::new(*duplicate_label, array_size), &array_size, - |b, _| { - b.iter(|| { - black_box( - udf.invoke_with_args(ScalarFunctionArgs { - args: vec![ColumnarValue::Array(array.clone())], - arg_fields: vec![ - Field::new("arr", array.data_type().clone(), false) - .into(), - ], - number_rows: NUM_ROWS, - return_field: Field::new( - "result", - array.data_type().clone(), - false, - ) - .into(), - config_options: Arc::new(ConfigOptions::default()), - }) - .unwrap(), - ) - }) - }, + |b, _| b.iter(|| invoke_unary_udf(&udf, &array, NUM_ROWS)), ); } } @@ -358,30 +360,80 @@ fn bench_array_distinct_sliced(c: &mut Criterion) { group.bench_with_input( BenchmarkId::from_parameter(array_size), &array_size, - |b, _| { - b.iter(|| { - black_box( - udf.invoke_with_args(ScalarFunctionArgs { - args: vec![ColumnarValue::Array(array.clone())], - arg_fields: vec![ - Field::new("arr", array.data_type().clone(), false) - .into(), - ], - number_rows: NUM_ROWS, - return_field: Field::new( - "result", - array.data_type().clone(), - false, - ) - .into(), - config_options: Arc::new(ConfigOptions::default()), - }) - .unwrap(), - ) - }) - }, + |b, _| b.iter(|| invoke_unary_udf(&udf, &array, NUM_ROWS)), + ); + } + group.finish(); +} + +fn create_sliced_float_array(backing_values: usize, visible_values: usize) -> ArrayRef { + assert!(visible_values > 0 && visible_values <= backing_values); + + let values = Float64Array::from( + (0..backing_values) + .map(|i| if i.is_multiple_of(2) { -0.0 } else { 0.0 }) + .collect::>(), + ); + let left_padding = (backing_values - visible_values) / 2; + let offsets = vec![ + 0, + left_padding as i32, + (left_padding + visible_values) as i32, + backing_values as i32, + ]; + let array = ListArray::try_new( + Arc::new(Field::new("item", DataType::Float64, true)), + OffsetBuffer::new(offsets.into()), + Arc::new(values), + None, + ) + .unwrap(); + + Arc::new(array.slice(1, 1)) +} + +fn create_unsliced_float_array(value_count: usize) -> ArrayRef { + let values = Float64Array::from( + (0..value_count) + .map(|i| if i.is_multiple_of(2) { -0.0 } else { 0.0 }) + .collect::>(), + ); + Arc::new(ListArray::new( + Arc::new(Field::new("item", DataType::Float64, true)), + OffsetBuffer::new(vec![0, value_count as i32].into()), + Arc::new(values), + None, + )) +} + +/// Keep the visible list fixed at one `-0.0` and one `0.0` while increasing +/// the backing values buffer. This isolates work outside the logical slice. +fn bench_array_distinct_sliced_float(c: &mut Criterion) { + let mut group = c.benchmark_group("array_distinct_sliced_float"); + let udf = ArrayDistinct::new(); + + for &backing_values in SLICED_FLOAT_BACKING_VALUES { + let array = create_sliced_float_array(backing_values, 2); + group.bench_with_input( + BenchmarkId::new("backing_values", backing_values), + &backing_values, + |b, _| b.iter(|| invoke_unary_udf(&udf, &array, 1)), ); } + + for &visible_values in SLICED_FLOAT_VISIBLE_VALUES { + let array = create_sliced_float_array(1024 * 1024, visible_values); + group.bench_with_input( + BenchmarkId::new("visible_values", visible_values), + &visible_values, + |b, _| b.iter(|| invoke_unary_udf(&udf, &array, 1)), + ); + } + + let array = create_unsliced_float_array(1024); + group.bench_function("unsliced_values_1024", |b| { + b.iter(|| invoke_unary_udf(&udf, &array, 1)) + }); group.finish(); } diff --git a/datafusion/functions-nested/src/except.rs b/datafusion/functions-nested/src/except.rs index 737a3122bbba..fa3e77711687 100644 --- a/datafusion/functions-nested/src/except.rs +++ b/datafusion/functions-nested/src/except.rs @@ -169,26 +169,32 @@ fn general_except( ) -> Result> { let converter = RowConverter::new(vec![SortField::new(l.value_type())])?; - // Normalize -0.0 → +0.0 so RowConverter (IEEE 754 totalOrder) groups - // ±0 together for both the rhs lookup set and the lhs probe. - let l_values_norm = normalize_float_zero(l.values()); - let r_values_norm = normalize_float_zero(r.values()); - - // Only convert the visible portion of the values array. For sliced - // ListArrays, values() returns the full underlying array but only - // elements between the first and last offset are referenced. + // ListArray::values() returns the full underlying array for sliced lists. + // Slice first so normalization only scans and, when -0.0 is present, + // allocates for values referenced by the logical array. + // Normalization keeps SQL signed-zero equality when rows are encoded. let l_first = l.offsets()[0].as_usize(); let l_len = l.offsets()[l.len()].as_usize() - l_first; - let l_values = converter.convert_columns(&[l_values_norm.slice(l_first, l_len)])?; + let l_values_norm = if l_first == 0 && l_len == l.values().len() { + normalize_float_zero(l.values()) + } else { + normalize_float_zero(&l.values().slice(l_first, l_len)) + }; + let l_rows = converter.convert_columns(&[l_values_norm.slice(0, l_len)])?; let r_first = r.offsets()[0].as_usize(); let r_len = r.offsets()[r.len()].as_usize() - r_first; - let r_values = converter.convert_columns(&[r_values_norm.slice(r_first, r_len)])?; + let r_values_norm = if r_first == 0 && r_len == r.values().len() { + normalize_float_zero(r.values()) + } else { + normalize_float_zero(&r.values().slice(r_first, r_len)) + }; + let r_rows = converter.convert_columns(&[r_values_norm.slice(0, r_len)])?; let mut offsets = Vec::::with_capacity(l.len() + 1); offsets.push(OffsetSize::usize_as(0)); - let mut indices: Vec = Vec::with_capacity(l_values.num_rows()); + let mut indices: Vec = Vec::with_capacity(l_rows.num_rows()); let mut dedup = HashSet::new(); let nulls = NullBuffer::union(l.nulls(), r.nulls()); @@ -207,13 +213,13 @@ fn general_except( } for element_index in r_start.as_usize() - r_first..r_end.as_usize() - r_first { - let right_row = r_values.row(element_index); + let right_row = r_rows.row(element_index); dedup.insert(right_row); } for element_index in l_start.as_usize() - l_first..l_end.as_usize() - l_first { - let left_row = l_values.row(element_index); + let left_row = l_rows.row(element_index); if dedup.insert(left_row) { - indices.push(element_index + l_first); + indices.push(element_index); } } @@ -245,9 +251,9 @@ fn general_except( #[cfg(test)] mod tests { - use super::ArrayExcept; - use arrow::array::{Array, AsArray, Int32Array, ListArray}; - use arrow::datatypes::{Field, Int32Type}; + use super::{ArrayExcept, general_except}; + use arrow::array::{Array, AsArray, Int32Array, LargeListArray, ListArray}; + use arrow::datatypes::{DataType, Field, Float64Type, Int32Type}; use datafusion_common::{Result, config::ConfigOptions}; use datafusion_expr::{ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl}; use std::sync::Arc; @@ -302,4 +308,30 @@ mod tests { Ok(()) } + + #[test] + fn test_array_except_sliced_float_large_lists() -> Result<()> { + let l = LargeListArray::from_iter_primitive::(vec![ + Some(vec![Some(99.0)]), + Some(vec![Some(-0.0), Some(0.0), Some(1.0), None, None]), + Some(vec![Some(-99.0)]), + ]) + .slice(1, 1); + let r = LargeListArray::from_iter_primitive::(vec![ + Some(vec![Some(98.0)]), + Some(vec![Some(0.0), None]), + Some(vec![Some(-98.0)]), + ]) + .slice(1, 1); + let DataType::LargeList(field) = l.data_type() else { + unreachable!() + }; + + let result = general_except::(&l, &r, field)?; + let values = result.value(0); + let values = values.as_primitive::(); + assert_eq!(values.len(), 1); + assert_eq!(values.value(0), 1.0); + Ok(()) + } } diff --git a/datafusion/functions-nested/src/set_ops.rs b/datafusion/functions-nested/src/set_ops.rs index 2214d3d35bb7..f2e98141b33b 100644 --- a/datafusion/functions-nested/src/set_ops.rs +++ b/datafusion/functions-nested/src/set_ops.rs @@ -351,29 +351,32 @@ fn generic_set_lists( let converter = RowConverter::new(vec![SortField::new(l.value_type())])?; - // Normalize -0.0 → +0.0 so RowConverter (which uses IEEE 754 totalOrder - // and treats ±0 as distinct) groups them together. Use the normalized - // arrays for both row conversion and the final output values. - let l_values_norm = normalize_float_zero(l.values()); - let r_values_norm = normalize_float_zero(r.values()); - - // Only convert the visible portion of the values array. For sliced - // ListArrays, values() returns the full underlying array but only - // elements between the first and last offset are referenced. + // ListArray::values() returns the full underlying array for sliced lists. + // Slice first so normalization only scans and, when -0.0 is present, + // allocates for values referenced by the logical array. + // Normalization keeps SQL signed-zero equality when rows are encoded. let l_first = l.offsets()[0].as_usize(); let l_len = l.offsets()[l.len()].as_usize() - l_first; - let l_values = l_values_norm.slice(l_first, l_len); - let rows_l = converter.convert_columns(&[Arc::clone(&l_values)])?; + let l_values_norm = if l_first == 0 && l_len == l.values().len() { + normalize_float_zero(l.values()) + } else { + normalize_float_zero(&l.values().slice(l_first, l_len)) + }; + let rows_l = converter.convert_columns(&[l_values_norm.slice(0, l_len)])?; let r_first = r.offsets()[0].as_usize(); let r_len = r.offsets()[r.len()].as_usize() - r_first; - let r_values = r_values_norm.slice(r_first, r_len); - let rows_r = converter.convert_columns(&[Arc::clone(&r_values)])?; + let r_values_norm = if r_first == 0 && r_len == r.values().len() { + normalize_float_zero(r.values()) + } else { + normalize_float_zero(&r.values().slice(r_first, r_len)) + }; + let rows_r = converter.convert_columns(&[r_values_norm.slice(0, r_len)])?; // Indices from the row converter are 0-based in the per-side slice; // concatenating those same slices lets indices map directly into the // combined values array. - let combined_values = concat(&[l_values.as_ref(), r_values.as_ref()])?; + let combined_values = concat(&[l_values_norm.as_ref(), r_values_norm.as_ref()])?; let r_offset = l_len; match set_op { @@ -565,18 +568,18 @@ fn general_array_distinct( let converter = RowConverter::new(vec![SortField::new(dt.clone())])?; - // Normalize -0.0 → +0.0 so RowConverter (which uses IEEE 754 totalOrder - // and treats ±0 as distinct) groups them together, and so the output - // carries the canonical sign. - let values_norm = normalize_float_zero(array.values()); - - // Only convert the visible portion of the values array. For sliced - // ListArrays, values() returns the full underlying array but only - // elements between the first and last offset are referenced. + // ListArray::values() returns the full underlying array for sliced lists. + // Slice first so normalization only scans and, when -0.0 is present, + // allocates for values referenced by the logical array. + // Normalization keeps SQL signed-zero equality and canonicalizes output. let first_offset = value_offsets[0].as_usize(); let visible_len = value_offsets[array.len()].as_usize() - first_offset; - let rows = - converter.convert_columns(&[values_norm.slice(first_offset, visible_len)])?; + let values_norm = if first_offset == 0 && visible_len == array.values().len() { + normalize_float_zero(array.values()) + } else { + normalize_float_zero(&array.values().slice(first_offset, visible_len)) + }; + let rows = converter.convert_columns(&[values_norm.slice(0, visible_len)])?; let mut indices: Vec = Vec::with_capacity(rows.num_rows()); let mut seen = HashSet::new(); @@ -598,15 +601,14 @@ fn general_array_distinct( for idx in start..end { let row = rows.row(idx); if seen.insert(row) { - indices.push(idx + first_offset); + indices.push(idx); } } offsets.push(last_offset + OffsetSize::usize_as(seen.len())); } // Gather distinct values in a single pass, using the computed `indices`. - // Indices are absolute positions in the (normalized) values array, so we - // can take directly from the full values. + // Indices are relative to the visible, normalized values array. // Use UInt64Array for LargeList to support values arrays exceeding u32::MAX. let final_values = if indices.is_empty() { new_empty_array(&dt) @@ -634,14 +636,17 @@ mod tests { use std::sync::Arc; use arrow::{ - array::{Array, AsArray, Int32Array, ListArray}, + array::{Array, AsArray, Int32Array, LargeListArray, ListArray}, buffer::OffsetBuffer, - datatypes::{DataType, Field, Int32Type}, + datatypes::{DataType, Field, Float64Type, Int32Type}, }; use datafusion_common::{DataFusionError, Result, config::ConfigOptions}; use datafusion_expr::{ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl}; - use crate::set_ops::{ArrayDistinct, ArrayIntersect, ArrayUnion, array_distinct_udf}; + use crate::set_ops::{ + ArrayDistinct, ArrayIntersect, ArrayUnion, SetOp, array_distinct_udf, + general_array_distinct, generic_set_lists, + }; /// Build two sliced ListArrays and return them along with the shared list /// field. @@ -678,6 +683,17 @@ mod tests { .collect() } + fn collect_f64_bits( + list: &arrow::array::GenericListArray, + row: usize, + ) -> Vec> { + list.value(row) + .as_primitive::() + .iter() + .map(|value| value.map(f64::to_bits)) + .collect() + } + #[test] fn test_array_union_sliced_lists() -> Result<()> { let (l, r, field) = make_sliced_pair(); @@ -761,6 +777,113 @@ mod tests { Ok(()) } + #[test] + fn test_sliced_float_set_ops_preserve_semantics() -> Result<()> { + let nan = f64::from_bits(0x7ff8_0000_0000_0042); + let l = ListArray::from_iter_primitive::(vec![ + Some(vec![Some(99.0)]), + Some(vec![ + Some(-0.0), + Some(0.0), + Some(1.0), + Some(nan), + Some(nan), + None, + None, + ]), + None, + Some(vec![Some(-99.0)]), + ]) + .slice(1, 2); + let r = ListArray::from_iter_primitive::(vec![ + Some(vec![Some(98.0)]), + Some(vec![Some(0.0), Some(2.0), Some(nan), None]), + Some(vec![Some(1.0)]), + Some(vec![Some(-98.0)]), + ]) + .slice(1, 2); + let DataType::List(field) = l.data_type() else { + unreachable!() + }; + + let distinct = general_array_distinct::(&l, field)?; + let distinct = distinct.as_list::(); + assert_eq!( + collect_f64_bits(distinct, 0), + vec![ + Some(0.0_f64.to_bits()), + Some(1.0_f64.to_bits()), + Some(nan.to_bits()), + None, + ] + ); + assert!(distinct.is_null(1)); + + let union = generic_set_lists::(&l, &r, Arc::clone(field), SetOp::Union)?; + let union = union.as_list::(); + assert_eq!( + collect_f64_bits(union, 0), + vec![ + Some(0.0_f64.to_bits()), + Some(1.0_f64.to_bits()), + Some(nan.to_bits()), + None, + Some(2.0_f64.to_bits()), + ] + ); + assert!(union.is_null(1)); + + let intersect = + generic_set_lists::(&l, &r, Arc::clone(field), SetOp::Intersect)?; + let intersect = intersect.as_list::(); + assert_eq!( + collect_f64_bits(intersect, 0), + vec![Some(0.0_f64.to_bits()), Some(nan.to_bits()), None] + ); + assert!(intersect.is_null(1)); + Ok(()) + } + + #[test] + fn test_array_distinct_sliced_float_large_list() -> Result<()> { + let list = LargeListArray::from_iter_primitive::(vec![ + Some(vec![Some(99.0)]), + Some(vec![Some(-0.0), Some(0.0), Some(1.0), None, None]), + Some(vec![Some(-99.0)]), + ]); + let sliced = list.slice(1, 1); + let DataType::LargeList(field) = sliced.data_type() else { + unreachable!() + }; + + let result = general_array_distinct::(&sliced, field)?; + let result = result.as_list::(); + assert_eq!( + collect_f64_bits(result, 0), + vec![Some(0.0_f64.to_bits()), Some(1.0_f64.to_bits()), None] + ); + Ok(()) + } + + #[test] + fn test_array_distinct_sliced_empty_float_list() -> Result<()> { + let list = ListArray::from_iter_primitive::(vec![ + Some(vec![Some(-0.0), Some(99.0)]), + Some(Vec::>::new()), + Some(vec![Some(-0.0), Some(0.0)]), + ]); + let sliced = list.slice(1, 1); + let DataType::List(field) = sliced.data_type() else { + unreachable!() + }; + + let result = general_array_distinct::(&sliced, field)?; + let result = result.as_list::(); + assert_eq!(result.len(), 1); + assert_eq!(result.value_length(0), 0); + Ok(()) + } + #[test] fn test_array_distinct_inner_nullability_result_type_match_return_type() -> Result<(), DataFusionError> {