Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 78 additions & 22 deletions datafusion/core/tests/fuzz_cases/equivalence/ordering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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]
Expand All @@ -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)?,
Expand All @@ -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));
Expand All @@ -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.
Expand All @@ -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)?;
Expand All @@ -118,6 +137,31 @@ fn test_ordering_satisfy_with_equivalence_complex_random() -> Result<()> {
Operator::Plus,
col("b", &test_schema)?,
)) as Arc<dyn PhysicalExpr>;
let binary = |lhs: Arc<dyn PhysicalExpr>, op, rhs: Arc<dyn PhysicalExpr>| {
Arc::new(BinaryExpr::new(lhs, op, rhs)) as Arc<dyn PhysicalExpr>
};
// 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)?,
Expand All @@ -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));
Expand All @@ -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(),
);
Expand Down Expand Up @@ -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![
Expand Down Expand Up @@ -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![
Expand Down
48 changes: 34 additions & 14 deletions datafusion/core/tests/fuzz_cases/equivalence/projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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]
Expand All @@ -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)?;
Expand Down Expand Up @@ -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).
Expand All @@ -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)?;
Expand All @@ -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
Expand All @@ -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)
});
Expand All @@ -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
Expand All @@ -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(),
Expand Down
16 changes: 10 additions & 6 deletions datafusion/core/tests/fuzz_cases/equivalence/properties.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};

Expand All @@ -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)?;
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading