From df3211c9a778a2772b74d86b95cf3bd7d9e367a8 Mon Sep 17 00:00:00 2001 From: Kevin-Li-2025 Date: Tue, 15 Sep 2026 00:30:43 +0400 Subject: [PATCH 1/3] fix: normalize only visible list values --- .../functions-nested/benches/array_set_ops.rs | 142 +++++++++----- datafusion/functions-nested/src/except.rs | 66 +++++-- datafusion/functions-nested/src/set_ops.rs | 183 +++++++++++++++--- 3 files changed, 299 insertions(+), 92 deletions(-) diff --git a/datafusion/functions-nested/benches/array_set_ops.rs b/datafusion/functions-nested/benches/array_set_ops.rs index d43bbdb577d06..fa5f26531aa63 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 737a3122bbbae..fa3e777116877 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 2214d3d35bb7b..f2e98141b33b6 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> { From 284a8cf65d5bfd5251e7ea5261d6c2d6af40728a Mon Sep 17 00:00:00 2001 From: Kevin-Li-2025 Date: Thu, 17 Sep 2026 18:58:25 +0400 Subject: [PATCH 2/3] address review --- .../functions-nested/benches/array_set_ops.rs | 45 ++-------------- datafusion/functions-nested/src/except.rs | 27 +++------- datafusion/functions-nested/src/set_ops.rs | 53 +++++++++---------- 3 files changed, 36 insertions(+), 89 deletions(-) diff --git a/datafusion/functions-nested/benches/array_set_ops.rs b/datafusion/functions-nested/benches/array_set_ops.rs index fa5f26531aa63..9c51974cd363c 100644 --- a/datafusion/functions-nested/benches/array_set_ops.rs +++ b/datafusion/functions-nested/benches/array_set_ops.rs @@ -38,11 +38,6 @@ 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); @@ -392,46 +387,14 @@ fn create_sliced_float_array(backing_values: usize, visible_values: usize) -> Ar 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. +/// Keep the visible list fixed at one `-0.0` and one `0.0` inside a much larger +/// backing values buffer to catch regressions that process values outside the slice. fn bench_array_distinct_sliced_float(c: &mut Criterion) { let mut group = c.benchmark_group("array_distinct_sliced_float"); let udf = ArrayDistinct::new(); + let array = create_sliced_float_array(1024 * 1024, 2); - 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| { + group.bench_function("backing_values_1048576_visible_values_2", |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 fa3e777116877..fd2b6a08977ed 100644 --- a/datafusion/functions-nested/src/except.rs +++ b/datafusion/functions-nested/src/except.rs @@ -17,6 +17,7 @@ //! [`ScalarUDFImpl`] definition for array_except function. +use crate::set_ops::normalize_visible_values; use crate::utils::{check_datatypes, make_scalar_function}; use arrow::array::new_null_array; use arrow::array::{ @@ -27,7 +28,7 @@ use arrow::buffer::{NullBuffer, OffsetBuffer}; use arrow::compute::take; use arrow::datatypes::{DataType, FieldRef}; use arrow::row::{RowConverter, SortField}; -use datafusion_common::utils::{ListCoercion, normalize_float_zero, take_function_args}; +use datafusion_common::utils::{ListCoercion, take_function_args}; use datafusion_common::{HashSet, Result, internal_err}; use datafusion_expr::{ ColumnarValue, Documentation, ScalarFunctionArgs, ScalarUDFImpl, Signature, @@ -169,27 +170,15 @@ fn general_except( ) -> Result> { let converter = RowConverter::new(vec![SortField::new(l.value_type())])?; - // 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. + // 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_first = l.offsets()[0].as_usize(); - let l_len = l.offsets()[l.len()].as_usize() - l_first; - 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 l_values_norm = normalize_visible_values(l); + let l_rows = converter.convert_columns(&[Arc::clone(&l_values_norm)])?; let r_first = r.offsets()[0].as_usize(); - let r_len = r.offsets()[r.len()].as_usize() - r_first; - 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 r_values_norm = normalize_visible_values(r); + let r_rows = converter.convert_columns(&[Arc::clone(&r_values_norm)])?; let mut offsets = Vec::::with_capacity(l.len() + 1); offsets.push(OffsetSize::usize_as(0)); diff --git a/datafusion/functions-nested/src/set_ops.rs b/datafusion/functions-nested/src/set_ops.rs index f2e98141b33b6..a847695a4bd0a 100644 --- a/datafusion/functions-nested/src/set_ops.rs +++ b/datafusion/functions-nested/src/set_ops.rs @@ -329,6 +329,18 @@ impl Display for SetOp { } } +pub(crate) fn normalize_visible_values( + array: &GenericListArray, +) -> ArrayRef { + let first = array.offsets()[0].as_usize(); + let len = array.offsets()[array.len()].as_usize() - first; + if first == 0 && len == array.values().len() { + normalize_float_zero(array.values()) + } else { + normalize_float_zero(&array.values().slice(first, len)) + } +} + fn generic_set_lists( l: &GenericListArray, r: &GenericListArray, @@ -351,27 +363,16 @@ fn generic_set_lists( let converter = RowConverter::new(vec![SortField::new(l.value_type())])?; - // 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. + // 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_first = l.offsets()[0].as_usize(); let l_len = l.offsets()[l.len()].as_usize() - l_first; - 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 l_values_norm = normalize_visible_values(l); + let rows_l = converter.convert_columns(&[Arc::clone(&l_values_norm)])?; - let r_first = r.offsets()[0].as_usize(); - let r_len = r.offsets()[r.len()].as_usize() - r_first; - 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)])?; + let r_values_norm = normalize_visible_values(r); + let rows_r = converter.convert_columns(&[Arc::clone(&r_values_norm)])?; // Indices from the row converter are 0-based in the per-side slice; // concatenating those same slices lets indices map directly into the @@ -568,18 +569,12 @@ fn general_array_distinct( let converter = RowConverter::new(vec![SortField::new(dt.clone())])?; - // 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. + // 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 first_offset = value_offsets[0].as_usize(); - let visible_len = value_offsets[array.len()].as_usize() - first_offset; - 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 values_norm = normalize_visible_values(array); + let rows = converter.convert_columns(&[Arc::clone(&values_norm)])?; let mut indices: Vec = Vec::with_capacity(rows.num_rows()); let mut seen = HashSet::new(); From 800a042cd8230e58f9a776604a41d0a7cc1c6d2a Mon Sep 17 00:00:00 2001 From: Kevin-Li-2025 Date: Thu, 17 Sep 2026 21:21:47 +0400 Subject: [PATCH 3/3] cleanup set operations --- datafusion/functions-nested/src/set_ops.rs | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/datafusion/functions-nested/src/set_ops.rs b/datafusion/functions-nested/src/set_ops.rs index a847695a4bd0a..5cd4819caefd7 100644 --- a/datafusion/functions-nested/src/set_ops.rs +++ b/datafusion/functions-nested/src/set_ops.rs @@ -366,19 +366,17 @@ fn generic_set_lists( // 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_first = l.offsets()[0].as_usize(); - let l_len = l.offsets()[l.len()].as_usize() - l_first; let l_values_norm = normalize_visible_values(l); - let rows_l = converter.convert_columns(&[Arc::clone(&l_values_norm)])?; - let r_values_norm = normalize_visible_values(r); + + let rows_l = converter.convert_columns(&[Arc::clone(&l_values_norm)])?; let rows_r = converter.convert_columns(&[Arc::clone(&r_values_norm)])?; // 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_norm.as_ref(), r_values_norm.as_ref()])?; - let r_offset = l_len; + let r_offset = l_values_norm.len(); match set_op { SetOp::Union => generic_set_loop::(