diff --git a/datafusion/core/tests/fuzz_cases/equivalence/ordering.rs b/datafusion/core/tests/fuzz_cases/equivalence/ordering.rs index 60b09976355e9..3864afd8a0cca 100644 --- a/datafusion/core/tests/fuzz_cases/equivalence/ordering.rs +++ b/datafusion/core/tests/fuzz_cases/equivalence/ordering.rs @@ -16,9 +16,10 @@ // under the License. use crate::fuzz_cases::equivalence::utils::{ - TestScalarUDF, contains_overflowable_arithmetic, create_random_schema, - create_test_params, create_test_schema_2, generate_table_for_eq_properties, - generate_table_for_orderings, is_table_same_after_sort, + NULL_PCTS, TestScalarUDF, assert_random_ordering_satisfy_is_sound, + contains_conservative_ordering_op, create_random_schema, create_test_params, + create_test_schema_2, generate_table_for_eq_properties, generate_table_for_orderings, + is_table_same_after_sort, }; use arrow::compute::SortOptions; use datafusion_common::Result; @@ -32,6 +33,8 @@ use datafusion_physical_expr::expressions::{BinaryExpr, col}; use datafusion_physical_expr_common::physical_expr::PhysicalExpr; use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr}; use itertools::Itertools; +use rand::SeedableRng; +use rand::rngs::StdRng; use std::sync::Arc; #[test] @@ -44,12 +47,16 @@ fn test_ordering_satisfy_with_equivalence_random() -> Result<()> { nulls_first: false, }; - for seed in 0..N_RANDOM_SCHEMA { + for (seed, &null_pct) in (0..N_RANDOM_SCHEMA).cartesian_product(NULL_PCTS) { // Create a random schema with random properties - let (test_schema, eq_properties) = create_random_schema(seed as u64)?; + let (test_schema, eq_properties) = create_random_schema(seed as u64, null_pct)?; // Generate a data that satisfies properties given - let table_data_with_properties = - generate_table_for_eq_properties(&eq_properties, N_ELEMENTS, N_DISTINCT)?; + let table_data_with_properties = generate_table_for_eq_properties( + &eq_properties, + N_ELEMENTS, + N_DISTINCT, + null_pct, + )?; let col_exprs = [ col("a", &test_schema)?, col("b", &test_schema)?, @@ -59,8 +66,16 @@ fn test_ordering_satisfy_with_equivalence_random() -> Result<()> { col("f", &test_schema)?, ]; + let mut options_rng = StdRng::seed_from_u64(seed as u64); for n_req in 1..=col_exprs.len() { for exprs in col_exprs.iter().combinations(n_req) { + assert_random_ordering_satisfy_is_sound( + &eq_properties, + &exprs, + &table_data_with_properties, + &mut options_rng, + &format!("seed: {seed}, null_pct: {null_pct}"), + )?; let sort_exprs = exprs .into_iter() .map(|expr| PhysicalSortExpr::new(Arc::clone(expr), SORT_OPTIONS)); @@ -72,7 +87,7 @@ fn test_ordering_satisfy_with_equivalence_random() -> Result<()> { &table_data_with_properties, )?; let err_msg = format!( - "Error in test case requirement:{ordering:?}, expected: {expected:?}, eq_properties {eq_properties}" + "Error in test case seed: {seed}, null_pct: {null_pct}, requirement:{ordering:?}, expected: {expected:?}, eq_properties {eq_properties}" ); // Check whether ordering_satisfy API result and // experimental result matches. @@ -98,12 +113,16 @@ fn test_ordering_satisfy_with_equivalence_complex_random() -> Result<()> { nulls_first: false, }; - for seed in 0..N_RANDOM_SCHEMA { + for (seed, &null_pct) in (0..N_RANDOM_SCHEMA).cartesian_product(NULL_PCTS) { // Create a random schema with random properties - let (test_schema, eq_properties) = create_random_schema(seed as u64)?; + let (test_schema, eq_properties) = create_random_schema(seed as u64, null_pct)?; // Generate a data that satisfies properties given - let table_data_with_properties = - generate_table_for_eq_properties(&eq_properties, N_ELEMENTS, N_DISTINCT)?; + let table_data_with_properties = generate_table_for_eq_properties( + &eq_properties, + N_ELEMENTS, + N_DISTINCT, + null_pct, + )?; let test_fun = Arc::new(ScalarUDF::new_from_impl(TestScalarUDF::new())); let col_a = col("a", &test_schema)?; @@ -118,6 +137,31 @@ fn test_ordering_satisfy_with_equivalence_complex_random() -> Result<()> { Operator::Plus, col("b", &test_schema)?, )) as Arc; + let binary = |lhs: Arc, op, rhs: Arc| { + Arc::new(BinaryExpr::new(lhs, op, rhs)) as Arc + }; + // Ordered when `a` and `b` lead orderings in opposite directions. + let a_gt_b = binary( + col("a", &test_schema)?, + Operator::Gt, + col("b", &test_schema)?, + ); + // `e` is constant, so each operand has the ordering of its column. + // `a` and `b` lead separate orderings whose NULLs sit in different + // rows, which is where Kleene `AND`/`OR` can break an ordering. + let a_gt_e = binary( + col("a", &test_schema)?, + Operator::Gt, + col("e", &test_schema)?, + ); + let b_gt_e = binary( + col("b", &test_schema)?, + Operator::Gt, + col("e", &test_schema)?, + ); + let a_gt_e_and_b_gt_e = + binary(Arc::clone(&a_gt_e), Operator::And, Arc::clone(&b_gt_e)); + let a_gt_e_or_b_gt_e = binary(a_gt_e, Operator::Or, b_gt_e); let exprs = [ col("a", &test_schema)?, col("b", &test_schema)?, @@ -127,10 +171,21 @@ fn test_ordering_satisfy_with_equivalence_complex_random() -> Result<()> { col("f", &test_schema)?, floor_a, a_plus_b, + a_gt_b, + a_gt_e_and_b_gt_e, + a_gt_e_or_b_gt_e, ]; + let mut options_rng = StdRng::seed_from_u64(seed as u64); for n_req in 1..=exprs.len() { for exprs in exprs.iter().combinations(n_req) { + assert_random_ordering_satisfy_is_sound( + &eq_properties, + &exprs, + &table_data_with_properties, + &mut options_rng, + &format!("seed: {seed}, null_pct: {null_pct}"), + )?; let sort_exprs = exprs .into_iter() .map(|expr| PhysicalSortExpr::new(Arc::clone(expr), SORT_OPTIONS)); @@ -142,19 +197,20 @@ fn test_ordering_satisfy_with_equivalence_complex_random() -> Result<()> { &table_data_with_properties, )?; let err_msg = format!( - "Error in test case requirement:{ordering:?}, expected: {expected:?}, eq_properties: {eq_properties}", + "Error in test case seed: {seed}, null_pct: {null_pct}, requirement:{ordering:?}, expected: {expected:?}, eq_properties: {eq_properties}", ); - // A rejection turns inconclusive only from the first `+`/`-` - // key onwards, since possible overflow makes an ordering - // underivable even when the sample happens to be sorted. A - // table sorted by the full ordering is sorted by every prefix - // of it, so a rejected arithmetic-free prefix still proves - // the rejection is genuine. + // A rejection turns inconclusive only from the first key with + // a conservative ordering rule onwards (see + // `contains_conservative_ordering_op`), since such a rule can + // make an ordering underivable even when the sample happens + // to be sorted. A table sorted by the full ordering is sorted + // by every prefix of it, so a rejected prefix without such + // keys still proves the rejection is genuine. let conclusive_prefix = LexOrdering::new( ordering .iter() .take_while(|sort_expr| { - !contains_overflowable_arithmetic(&sort_expr.expr) + !contains_conservative_ordering_op(&sort_expr.expr) }) .cloned(), ); @@ -197,7 +253,7 @@ fn test_ordering_satisfy_with_equivalence() -> Result<()> { nulls_first: true, }; let table_data_with_properties = - generate_table_for_eq_properties(&eq_properties, 625, 5)?; + generate_table_for_eq_properties(&eq_properties, 625, 5, 0.0)?; // First element in the tuple stores vector of requirement, second element is the expected return value for ordering_satisfy function let requirements = vec![ @@ -372,7 +428,7 @@ fn test_ordering_satisfy_on_data() -> Result<()> { ]; let orderings = convert_to_orderings(&orderings); - let batch = generate_table_for_orderings(orderings, schema, 1000, 10)?; + let batch = generate_table_for_orderings(orderings, schema, 1000, 10, 0.0)?; // [a ASC, c ASC, d ASC] cannot be deduced let ordering = vec![ diff --git a/datafusion/core/tests/fuzz_cases/equivalence/projection.rs b/datafusion/core/tests/fuzz_cases/equivalence/projection.rs index 9593e1cf11565..c71dc482ec2cd 100644 --- a/datafusion/core/tests/fuzz_cases/equivalence/projection.rs +++ b/datafusion/core/tests/fuzz_cases/equivalence/projection.rs @@ -16,8 +16,9 @@ // under the License. use crate::fuzz_cases::equivalence::utils::{ - TestScalarUDF, apply_projection, contains_overflowable_arithmetic, - create_random_schema, generate_table_for_eq_properties, is_table_same_after_sort, + NULL_PCTS, TestScalarUDF, apply_projection, assert_random_ordering_satisfy_is_sound, + contains_conservative_ordering_op, create_random_schema, + generate_table_for_eq_properties, is_table_same_after_sort, }; use arrow::compute::SortOptions; use datafusion_common::Result; @@ -29,6 +30,8 @@ use datafusion_physical_expr::{PhysicalExprRef, ScalarFunctionExpr}; use datafusion_physical_expr_common::physical_expr::PhysicalExpr; use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr}; use itertools::Itertools; +use rand::SeedableRng; +use rand::rngs::StdRng; use std::sync::Arc; #[test] @@ -37,12 +40,16 @@ fn project_orderings_random() -> Result<()> { const N_ELEMENTS: usize = 125; const N_DISTINCT: usize = 5; - for seed in 0..N_RANDOM_SCHEMA { + for (seed, &null_pct) in (0..N_RANDOM_SCHEMA).cartesian_product(NULL_PCTS) { // Create a random schema with random properties - let (test_schema, eq_properties) = create_random_schema(seed as u64)?; + let (test_schema, eq_properties) = create_random_schema(seed as u64, null_pct)?; // Generate a data that satisfies properties given - let table_data_with_properties = - generate_table_for_eq_properties(&eq_properties, N_ELEMENTS, N_DISTINCT)?; + let table_data_with_properties = generate_table_for_eq_properties( + &eq_properties, + N_ELEMENTS, + N_DISTINCT, + null_pct, + )?; // Floor(a) let test_fun = Arc::new(ScalarUDF::new_from_impl(TestScalarUDF::new())); let col_a = col("a", &test_schema)?; @@ -84,7 +91,7 @@ fn project_orderings_random() -> Result<()> { // Make sure each ordering after projection is valid. for ordering in projected_eq.oeq_class().iter() { let err_msg = format!( - "Error in test case ordering:{ordering:?}, eq_properties {eq_properties}, proj_exprs: {proj_exprs:?}", + "Error in test case seed: {seed}, null_pct: {null_pct}, ordering:{ordering:?}, eq_properties {eq_properties}, proj_exprs: {proj_exprs:?}", ); // Since ordered section satisfies schema, we expect // that result will be same after sort (e.g sort was unnecessary). @@ -111,12 +118,16 @@ fn ordering_satisfy_after_projection_random() -> Result<()> { nulls_first: false, }; - for seed in 0..N_RANDOM_SCHEMA { + for (seed, &null_pct) in (0..N_RANDOM_SCHEMA).cartesian_product(NULL_PCTS) { // Create a random schema with random properties - let (test_schema, eq_properties) = create_random_schema(seed as u64)?; + let (test_schema, eq_properties) = create_random_schema(seed as u64, null_pct)?; // Generate a data that satisfies properties given - let table_data_with_properties = - generate_table_for_eq_properties(&eq_properties, N_ELEMENTS, N_DISTINCT)?; + let table_data_with_properties = generate_table_for_eq_properties( + &eq_properties, + N_ELEMENTS, + N_DISTINCT, + null_pct, + )?; // Floor(a) let test_fun = Arc::new(ScalarUDF::new_from_impl(TestScalarUDF::new())); let col_a = col("a", &test_schema)?; @@ -143,6 +154,7 @@ fn ordering_satisfy_after_projection_random() -> Result<()> { (a_plus_b, "a+b"), ]; + let mut options_rng = StdRng::seed_from_u64(seed as u64); for n_req in 0..=proj_exprs.len() { for proj_exprs in proj_exprs.iter().combinations(n_req) { let proj_exprs = proj_exprs @@ -166,6 +178,13 @@ fn ordering_satisfy_after_projection_random() -> Result<()> { for n_req in 1..=projected_exprs.len() { for exprs in projected_exprs.iter().combinations(n_req) { + assert_random_ordering_satisfy_is_sound( + &projected_eq, + &exprs, + &projected_batch, + &mut options_rng, + &format!("seed: {seed}, null_pct: {null_pct}"), + )?; let sort_exprs = exprs.into_iter().map(|expr| { PhysicalSortExpr::new(Arc::clone(expr), SORT_OPTIONS) }); @@ -177,10 +196,11 @@ fn ordering_satisfy_after_projection_random() -> Result<()> { let expected = is_table_same_after_sort(ordering.clone(), &projected_batch)?; let err_msg = format!( - "Error in test case requirement:{ordering:?}, expected: {expected:?}, eq_properties: {eq_properties}, projected_eq: {projected_eq}, projection_mapping: {projection_mapping:?}" + "Error in test case seed: {seed}, null_pct: {null_pct}, requirement:{ordering:?}, expected: {expected:?}, eq_properties: {eq_properties}, projected_eq: {projected_eq}, projection_mapping: {projection_mapping:?}" ); // Same reasoning as in `ordering.rs`: only keys from - // the first `+`/`-` source onwards are inconclusive, + // the first source with a conservative ordering rule + // onwards are inconclusive, // so assert on the longest prefix without one. let conclusive_prefix = LexOrdering::new( ordering @@ -190,7 +210,7 @@ fn ordering_satisfy_after_projection_random() -> Result<()> { targets .iter() .any(|(target, _)| target.eq(&sort_expr.expr)) - && contains_overflowable_arithmetic(source) + && contains_conservative_ordering_op(source) }) }) .cloned(), diff --git a/datafusion/core/tests/fuzz_cases/equivalence/properties.rs b/datafusion/core/tests/fuzz_cases/equivalence/properties.rs index 1490eb08a0291..9f04457841eca 100644 --- a/datafusion/core/tests/fuzz_cases/equivalence/properties.rs +++ b/datafusion/core/tests/fuzz_cases/equivalence/properties.rs @@ -18,7 +18,7 @@ use std::sync::Arc; use crate::fuzz_cases::equivalence::utils::{ - TestScalarUDF, create_random_schema, generate_table_for_eq_properties, + NULL_PCTS, TestScalarUDF, create_random_schema, generate_table_for_eq_properties, is_table_same_after_sort, }; @@ -37,12 +37,16 @@ fn test_find_longest_permutation_random() -> Result<()> { const N_ELEMENTS: usize = 125; const N_DISTINCT: usize = 5; - for seed in 0..N_RANDOM_SCHEMA { + for (seed, &null_pct) in (0..N_RANDOM_SCHEMA).cartesian_product(NULL_PCTS) { // Create a random schema with random properties - let (test_schema, eq_properties) = create_random_schema(seed as u64)?; + let (test_schema, eq_properties) = create_random_schema(seed as u64, null_pct)?; // Generate a data that satisfies properties given - let table_data_with_properties = - generate_table_for_eq_properties(&eq_properties, N_ELEMENTS, N_DISTINCT)?; + let table_data_with_properties = generate_table_for_eq_properties( + &eq_properties, + N_ELEMENTS, + N_DISTINCT, + null_pct, + )?; let test_fun = Arc::new(ScalarUDF::new_from_impl(TestScalarUDF::new())); let col_a = col("a", &test_schema)?; @@ -88,7 +92,7 @@ fn test_find_longest_permutation_random() -> Result<()> { ); let err_msg = format!( - "Error in test case ordering:{ordering:?}, eq_properties: {eq_properties}" + "Error in test case seed: {seed}, null_pct: {null_pct}, ordering:{ordering:?}, eq_properties: {eq_properties}" ); assert_eq!(ordering.len(), indices.len(), "{err_msg}"); // Since ordered section satisfies schema, we expect diff --git a/datafusion/core/tests/fuzz_cases/equivalence/utils.rs b/datafusion/core/tests/fuzz_cases/equivalence/utils.rs index ca73db3ae99ec..611b64dc1274c 100644 --- a/datafusion/core/tests/fuzz_cases/equivalence/utils.rs +++ b/datafusion/core/tests/fuzz_cases/equivalence/utils.rs @@ -20,13 +20,14 @@ use std::sync::Arc; use arrow::array::{ArrayRef, Float32Array, Float64Array, RecordBatch, UInt32Array}; use arrow::compute::{SortColumn, SortOptions, lexsort_to_indices, take_record_batch}; -use arrow::datatypes::{DataType, Field, Schema, SchemaRef}; +use arrow::datatypes::{DataType, Field, FieldRef, Schema, SchemaRef}; use datafusion_common::tree_node::TreeNode; use datafusion_common::utils::{compare_rows, get_row_at_idx}; use datafusion_common::{Result, exec_err, internal_datafusion_err, plan_err}; use datafusion_expr::sort_properties::{ExprProperties, SortProperties}; use datafusion_expr::{ - ColumnarValue, Operator, ScalarFunctionArgs, ScalarUDFImpl, Signature, Volatility, + ColumnarValue, Operator, ReturnFieldArgs, ScalarFunctionArgs, ScalarUDFImpl, + Signature, Volatility, }; use datafusion_physical_expr::equivalence::{ EquivalenceClass, ProjectionMapping, convert_to_orderings, @@ -83,8 +84,21 @@ pub fn create_test_schema_2() -> Result { /// where /// Column [a=f] (e.g they are aliases). /// Column e is constant. -pub fn create_random_schema(seed: u64) -> Result<(SchemaRef, EquivalenceProperties)> { - let test_schema = create_test_schema_2()?; +/// +/// Columns are declared nullable only when `null_pct > 0.0`, so the schema +/// tells `ordering_satisfy` whether `nulls_first` can matter for the data. +pub fn create_random_schema( + seed: u64, + null_pct: f64, +) -> Result<(SchemaRef, EquivalenceProperties)> { + let nullable = null_pct > 0.0; + let test_schema = Arc::new(Schema::new( + create_test_schema_2()? + .fields() + .iter() + .map(|field| field.as_ref().clone().with_nullable(nullable)) + .collect::>(), + )); let col_a = &col("a", &test_schema)?; let col_b = &col("b", &test_schema)?; let col_c = &col("c", &test_schema)?; @@ -103,11 +117,6 @@ pub fn create_random_schema(seed: u64) -> Result<(SchemaRef, EquivalenceProperti let mut rng = StdRng::seed_from_u64(seed); let mut remaining_exprs = col_exprs[0..4].to_vec(); // only a, b, c, d are sorted - let options_asc = SortOptions { - descending: false, - nulls_first: false, - }; - while !remaining_exprs.is_empty() { let n_sort_expr = rng.random_range(1..remaining_exprs.len() + 1); remaining_exprs.shuffle(&mut rng); @@ -117,7 +126,7 @@ pub fn create_random_schema(seed: u64) -> Result<(SchemaRef, EquivalenceProperti .drain(0..n_sort_expr) .map(|expr| PhysicalSortExpr { expr: Arc::clone(expr), - options: options_asc, + options: random_sort_options(&mut rng), }); eq_properties.add_ordering(ordering); @@ -126,6 +135,14 @@ pub fn create_random_schema(seed: u64) -> Result<(SchemaRef, EquivalenceProperti Ok((test_schema, eq_properties)) } +/// Picks the direction and NULL placement of one sort key. +pub fn random_sort_options(rng: &mut StdRng) -> SortOptions { + SortOptions { + descending: rng.random(), + nulls_first: rng.random(), + } +} + // Apply projection to the input_data, return projected equivalence properties and record batch pub fn apply_projection( proj_exprs: impl IntoIterator, String)>, @@ -210,15 +227,29 @@ fn add_equal_conditions_test() -> Result<()> { Ok(()) } -/// Returns `true` if `expr` contains a `+` or `-` anywhere in its tree. +/// Returns `true` if `expr` contains, anywhere in its tree, an operator whose +/// ordering rule is deliberately conservative. /// -/// The equivalence framework conservatively discards orderings derived from -/// `+`/`-` expressions, because wrapping overflow can break them over the -/// type's full domain even when a finite batch happens to remain sorted. -pub fn contains_overflowable_arithmetic(expr: &Arc) -> bool { +/// The equivalence framework discards orderings in cases where a finite batch +/// can still happen to be sorted, so a rejection is not conclusive: +/// - `+`/`-`: wrapping overflow can break the ordering over the type's full +/// domain. +/// - comparisons, `AND`, `OR`: the rules keep an ordering only for specific +/// NULL placements, without checking whether the operands can be NULL. +pub fn contains_conservative_ordering_op(expr: &Arc) -> bool { expr.exists(|e| { Ok(e.downcast_ref::().is_some_and(|binary| { - matches!(binary.op(), Operator::Plus | Operator::Minus) + matches!( + binary.op(), + Operator::Plus + | Operator::Minus + | Operator::Gt + | Operator::GtEq + | Operator::Lt + | Operator::LtEq + | Operator::And + | Operator::Or + ) })) }) .unwrap() @@ -288,6 +319,34 @@ pub fn is_table_same_after_sort( Ok(sorted_indices == original_indices) } +/// Sorts `exprs` with random [`SortOptions`] and asserts that `eq_properties` +/// never claims an ordering that `batch` does not have. +/// +/// Only this direction is checked: a single batch cannot tell an underivable +/// ordering from one that is sorted by coincidence, e.g. when NULLs of +/// independently sorted columns end up in the same rows. +pub fn assert_random_ordering_satisfy_is_sound( + eq_properties: &EquivalenceProperties, + exprs: &[&Arc], + batch: &RecordBatch, + rng: &mut StdRng, + context: &str, +) -> Result<()> { + let sort_exprs = exprs + .iter() + .map(|expr| PhysicalSortExpr::new(Arc::clone(expr), random_sort_options(rng))); + let Some(ordering) = LexOrdering::new(sort_exprs) else { + unreachable!("Test should always produce non-degenerate orderings"); + }; + if eq_properties.ordering_satisfy(ordering.clone())? { + assert!( + is_table_same_after_sort(ordering.clone(), batch)?, + "{context}, random requirement: {ordering:?}, eq_properties: {eq_properties}" + ); + } + Ok(()) +} + // If we already generated a random result for one of the // expressions in the equivalence classes. For other expressions in the same // equivalence class use same result. This util gets already calculated result, when available. @@ -369,6 +428,7 @@ pub fn generate_table_for_eq_properties( eq_properties: &EquivalenceProperties, n_elem: usize, n_distinct: usize, + null_pct: f64, ) -> Result { let mut rng = StdRng::seed_from_u64(23); @@ -377,10 +437,7 @@ pub fn generate_table_for_eq_properties( // Utility closure to generate random array let mut generate_random_array = |num_elems: usize, max_val: usize| -> ArrayRef { - let values: Vec = (0..num_elems) - .map(|_| rng.random_range(0..max_val) as f64 / 2.0) - .collect(); - Arc::new(Float64Array::from_iter_values(values)) + generate_random_f64_array(num_elems, max_val, null_pct, &mut rng) }; // Fill constant columns @@ -450,6 +507,7 @@ pub fn generate_table_for_orderings( schema: SchemaRef, n_elem: usize, n_distinct: usize, + null_pct: f64, ) -> Result { let mut rng = StdRng::seed_from_u64(23); @@ -463,7 +521,7 @@ pub fn generate_table_for_orderings( .map(|field| { ( field.name(), - generate_random_f64_array(n_elem, n_distinct, &mut rng), + generate_random_f64_array(n_elem, n_distinct, null_pct, &mut rng), ) }) .collect::>(); @@ -503,16 +561,25 @@ pub fn generate_table_for_orderings( Ok(batch) } +pub const NULL_PCTS: &[f64] = &[0.0, 0.1, 0.5]; + // Utility function to generate random f64 array fn generate_random_f64_array( n_elems: usize, n_distinct: usize, + null_pct: f64, rng: &mut StdRng, ) -> ArrayRef { - let values: Vec = (0..n_elems) - .map(|_| rng.random_range(0..n_distinct) as f64 / 2.0) - .collect(); - Arc::new(Float64Array::from_iter_values(values)) + let values = (0..n_elems) + .map(|_| { + if rng.random::() < null_pct { + None + } else { + Some(rng.random_range(0..n_distinct) as f64 / 2.0) + } + }) + .collect::(); + Arc::new(values) } // Helper function to get sort columns from a batch @@ -562,6 +629,17 @@ impl ScalarUDFImpl for TestScalarUDF { } } + fn return_field_from_args(&self, args: ReturnFieldArgs) -> Result { + let arg_field = &args.arg_fields[0]; + let return_type = self.return_type(&[arg_field.data_type().clone()])?; + // floor maps NULL to NULL and never produces NULL otherwise. + Ok(Arc::new(Field::new( + self.name(), + return_type, + arg_field.is_nullable(), + ))) + } + fn output_ordering(&self, input: &[ExprProperties]) -> Result { Ok(input[0].sort_properties) }