From cfd0f594c206a649ee785e896c64743fc2c21630 Mon Sep 17 00:00:00 2001 From: AndreaBozzo Date: Sat, 18 Jul 2026 09:50:44 +0200 Subject: [PATCH] fix(datafusion): classify filter pushdown accurately --- .../src/physical_plan/expr_to_predicate.rs | 327 +++++++++++++----- .../integrations/datafusion/src/table/mod.rs | 13 +- .../tests/integration_datafusion_test.rs | 30 +- .../df_test/binary_predicate_pushdown.slt | 9 +- .../df_test/boolean_predicate_pushdown.slt | 27 +- .../slts/df_test/like_predicate_pushdown.slt | 9 +- .../df_test/timestamp_predicate_pushdown.slt | 33 +- 7 files changed, 308 insertions(+), 140 deletions(-) diff --git a/crates/integrations/datafusion/src/physical_plan/expr_to_predicate.rs b/crates/integrations/datafusion/src/physical_plan/expr_to_predicate.rs index fb5440a98a..d2fb9b10f3 100644 --- a/crates/integrations/datafusion/src/physical_plan/expr_to_predicate.rs +++ b/crates/integrations/datafusion/src/physical_plan/expr_to_predicate.rs @@ -19,19 +19,59 @@ use std::vec; use datafusion::arrow::datatypes::DataType; use datafusion::logical_expr::expr::ScalarFunction; -use datafusion::logical_expr::{BinaryExpr, Expr, Like, Operator}; +use datafusion::logical_expr::{BinaryExpr, Expr, Like, Operator, TableProviderFilterPushDown}; use datafusion::scalar::ScalarValue; use iceberg::expr::{BinaryExpression, Predicate, PredicateOperator, Reference, UnaryExpression}; use iceberg::spec::{Datum, PrimitiveLiteral}; -// A datafusion expression could be an Iceberg predicate, column, or literal. -enum TransformedResult { +// A DataFusion expression could be an Iceberg predicate, column, or literal. +enum TransformedValue { Predicate(Predicate), Column(Reference), Literal(Datum), NotTransformed, } +struct TransformedResult { + value: TransformedValue, + exact: bool, +} + +impl TransformedResult { + fn predicate(predicate: Predicate, exact: bool) -> Self { + Self { + value: TransformedValue::Predicate(predicate), + exact, + } + } + + fn column(column: Reference) -> Self { + Self { + value: TransformedValue::Column(column), + exact: true, + } + } + + fn literal(literal: Datum) -> Self { + Self { + value: TransformedValue::Literal(literal), + exact: true, + } + } + + fn not_transformed() -> Self { + Self { + value: TransformedValue::NotTransformed, + exact: false, + } + } + + fn mark_inexact(mut self) -> Self { + self.exact = false; + self + } +} + enum OpTransformedResult { Operator(PredicateOperator), And, @@ -45,23 +85,27 @@ enum OpTransformedResult { pub fn convert_filters_to_predicate(filters: &[Expr]) -> Option { filters .iter() - .filter_map(convert_filter_to_predicate) + .filter_map(|filter| convert_filter_to_predicate(filter).map(|(predicate, _)| predicate)) .reduce(Predicate::and) } -fn convert_filter_to_predicate(expr: &Expr) -> Option { - match to_iceberg_predicate(expr) { - TransformedResult::Predicate(predicate) => Some(predicate), - TransformedResult::Column(column) => { +fn convert_filter_to_predicate(expr: &Expr) -> Option<(Predicate, bool)> { + let transformed = to_iceberg_predicate(expr); + match transformed.value { + TransformedValue::Predicate(predicate) => Some((predicate, transformed.exact)), + TransformedValue::Column(column) => { // A bare column in a filter context represents a boolean column check // Convert it to: column = true - Some(Predicate::Binary(BinaryExpression::new( - PredicateOperator::Eq, - column, - Datum::bool(true), - ))) + Some(( + Predicate::Binary(BinaryExpression::new( + PredicateOperator::Eq, + column, + Datum::bool(true), + )), + transformed.exact, + )) } - TransformedResult::Literal(_) => { + TransformedValue::Literal(_) => { // Literal values in filter context cannot be pushed down None } @@ -69,6 +113,14 @@ fn convert_filter_to_predicate(expr: &Expr) -> Option { } } +pub(crate) fn classify_filter_pushdown(expr: &Expr) -> TableProviderFilterPushDown { + match convert_filter_to_predicate(expr) { + Some((_, true)) => TableProviderFilterPushDown::Exact, + Some((_, false)) => TableProviderFilterPushDown::Inexact, + None => TableProviderFilterPushDown::Unsupported, + } +} + fn to_iceberg_predicate(expr: &Expr) -> TransformedResult { match expr { Expr::BinaryExpr(binary) => { @@ -79,64 +131,82 @@ fn to_iceberg_predicate(expr: &Expr) -> TransformedResult { OpTransformedResult::Operator(op) => to_iceberg_binary_predicate(left, right, op), OpTransformedResult::And => to_iceberg_and_predicate(left, right), OpTransformedResult::Or => to_iceberg_or_predicate(left, right), - OpTransformedResult::NotTransformed => TransformedResult::NotTransformed, + OpTransformedResult::NotTransformed => TransformedResult::not_transformed(), } } Expr::Not(exp) => { let expr = to_iceberg_predicate(exp); - match expr { - TransformedResult::Predicate(p) => TransformedResult::Predicate(!p), - TransformedResult::Column(column) => { + match expr.value { + TransformedValue::Predicate(p) if expr.exact => TransformedResult::predicate( + !p, + !matches!( + exp.as_ref(), + Expr::ScalarFunction(ScalarFunction { func, .. }) if func.name() == "isnan" + ), + ), + TransformedValue::Column(column) => { // NOT of a bare boolean column: NOT col => col = false - TransformedResult::Predicate(Predicate::Binary(BinaryExpression::new( - PredicateOperator::Eq, - column, - Datum::bool(false), - ))) + TransformedResult::predicate( + Predicate::Binary(BinaryExpression::new( + PredicateOperator::Eq, + column, + Datum::bool(false), + )), + expr.exact, + ) } - _ => TransformedResult::NotTransformed, + // Negating an inexact predicate can turn a safe superset into + // a subset and incorrectly prune matching rows. + _ => TransformedResult::not_transformed(), } } - Expr::Column(column) => TransformedResult::Column(Reference::new(column.name())), + Expr::Column(column) => TransformedResult::column(Reference::new(column.name())), Expr::Literal(literal, _) => match scalar_value_to_datum(literal) { - Some(data) => TransformedResult::Literal(data), - None => TransformedResult::NotTransformed, + Some(data) => TransformedResult::literal(data), + None => TransformedResult::not_transformed(), }, Expr::InList(inlist) => { let mut datums = vec![]; + let mut exact = true; for expr in &inlist.list { let p = to_iceberg_predicate(expr); - match p { - TransformedResult::Literal(l) => datums.push(l), - _ => return TransformedResult::NotTransformed, + exact &= p.exact; + match p.value { + TransformedValue::Literal(l) => datums.push(l), + _ => return TransformedResult::not_transformed(), } } let expr = to_iceberg_predicate(&inlist.expr); - match expr { - TransformedResult::Column(r) => match inlist.negated { - false => TransformedResult::Predicate(r.is_in(datums)), - true => TransformedResult::Predicate(r.is_not_in(datums)), + exact &= expr.exact; + match expr.value { + TransformedValue::Column(r) => match inlist.negated { + false => TransformedResult::predicate(r.is_in(datums), exact), + // The Arrow reader conservatively matches NOT IN when the + // column is absent from an older data file. + true => TransformedResult::predicate(r.is_not_in(datums), false), }, - _ => TransformedResult::NotTransformed, + _ => TransformedResult::not_transformed(), } } Expr::IsNull(expr) => { let p = to_iceberg_predicate(expr); - match p { - TransformedResult::Column(r) => TransformedResult::Predicate(Predicate::Unary( - UnaryExpression::new(PredicateOperator::IsNull, r), - )), - _ => TransformedResult::NotTransformed, + match p.value { + TransformedValue::Column(r) => TransformedResult::predicate( + Predicate::Unary(UnaryExpression::new(PredicateOperator::IsNull, r)), + p.exact, + ), + _ => TransformedResult::not_transformed(), } } Expr::IsNotNull(expr) => { let p = to_iceberg_predicate(expr); - match p { - TransformedResult::Column(r) => TransformedResult::Predicate(Predicate::Unary( - UnaryExpression::new(PredicateOperator::NotNull, r), - )), - _ => TransformedResult::NotTransformed, + match p.value { + TransformedValue::Column(r) => TransformedResult::predicate( + Predicate::Unary(UnaryExpression::new(PredicateOperator::NotNull, r)), + p.exact, + ), + _ => TransformedResult::not_transformed(), } } Expr::Cast(c) => { @@ -144,9 +214,12 @@ fn to_iceberg_predicate(expr: &Expr) -> TransformedResult { { // Casts to date truncate the expression, we cannot simply extract it as it // can create erroneous predicates. - return TransformedResult::NotTransformed; + return TransformedResult::not_transformed(); } - to_iceberg_predicate(&c.expr) + // Keep using the cast-stripped predicate for pruning, but retain the + // original DataFusion filter because this layer cannot prove the cast + // preserves comparison semantics. + to_iceberg_predicate(&c.expr).mark_inexact() } Expr::Like(Like { negated, @@ -160,16 +233,19 @@ fn to_iceberg_predicate(expr: &Expr) -> TransformedResult { // push down case-insensitive LIKE (ILIKE) patterns // Escape characters are also not supported for pushdown if escape_char.is_some() || *case_insensitive { - return TransformedResult::NotTransformed; + return TransformedResult::not_transformed(); } // Extract the pattern string let pattern_str = match to_iceberg_predicate(pattern) { - TransformedResult::Literal(d) => match d.literal() { + TransformedResult { + value: TransformedValue::Literal(d), + .. + } => match d.literal() { PrimitiveLiteral::String(s) => s.clone(), - _ => return TransformedResult::NotTransformed, + _ => return TransformedResult::not_transformed(), }, - _ => return TransformedResult::NotTransformed, + _ => return TransformedResult::not_transformed(), }; // Check if it's a simple prefix pattern (ends with % and no other wildcards) @@ -181,8 +257,11 @@ fn to_iceberg_predicate(expr: &Expr) -> TransformedResult { // Get the column reference let column = match to_iceberg_predicate(expr) { - TransformedResult::Column(r) => r, - _ => return TransformedResult::NotTransformed, + TransformedResult { + value: TransformedValue::Column(r), + .. + } => r, + _ => return TransformedResult::not_transformed(), }; // Create the appropriate predicate @@ -192,16 +271,18 @@ fn to_iceberg_predicate(expr: &Expr) -> TransformedResult { column.starts_with(Datum::string(prefix)) }; - TransformedResult::Predicate(predicate) + // The Arrow reader conservatively matches NOT STARTS WITH when + // the column is absent from an older data file. + TransformedResult::predicate(predicate, !*negated) } else { // Complex LIKE patterns cannot be pushed down - TransformedResult::NotTransformed + TransformedResult::not_transformed() } } Expr::ScalarFunction(ScalarFunction { func, args }) => { scalar_function_to_iceberg_predicate(func.name(), args) } - _ => TransformedResult::NotTransformed, + _ => TransformedResult::not_transformed(), } } @@ -227,10 +308,10 @@ fn to_iceberg_operation(op: Operator) -> OpTransformedResult { fn scalar_function_to_iceberg_predicate(func_name: &str, args: &[Expr]) -> TransformedResult { match func_name { "isnan" if args.len() == 1 => match resolve_nan_preserving_reference(&args[0]) { - Some(r) => TransformedResult::Predicate(r.is_nan()), - None => TransformedResult::NotTransformed, + Some(r) => TransformedResult::predicate(r.is_nan(), true), + None => TransformedResult::not_transformed(), }, - _ => TransformedResult::NotTransformed, + _ => TransformedResult::not_transformed(), } } @@ -238,13 +319,11 @@ fn scalar_function_to_iceberg_predicate(func_name: &str, args: &[Expr]) -> Trans /// [`Reference`] such that `isnan(arg)` is logically equivalent to /// `isnan(reference)`. /// -/// Filter pushdown is reported as `Inexact` (see -/// [`IcebergTableProvider::supports_filters_pushdown`]), so DataFusion -/// re-applies the original predicate after scanning. We therefore only need the -/// pushed-down predicate to be implied by the original filter (it may match -/// extra rows, but must never drop a matching one). Every transformation handled -/// here preserves NaN-ness *exactly* — the result is NaN if and only if the -/// wrapped column is NaN — so both `isnan(arg)` and `NOT isnan(arg)` are sound: +/// Positive `isnan` filters produced here may be reported as `Exact` by +/// [`IcebergTableProvider::supports_filters_pushdown`], so every transformation +/// handled here must preserve NaN-ness in both directions: the result is NaN if +/// and only if the wrapped column is NaN. Negated `isnan` remains `Inexact` +/// because DataFusion and the Iceberg Arrow predicate differ for null values. /// /// * negation: `-x` is NaN iff `x` is NaN /// * `abs(x)`: `abs(x)` is NaN iff `x` is NaN @@ -359,22 +438,24 @@ fn to_iceberg_and_predicate( left: TransformedResult, right: TransformedResult, ) -> TransformedResult { - match (left, right) { - (TransformedResult::Predicate(left), TransformedResult::Predicate(right)) => { - TransformedResult::Predicate(left.and(right)) + let exact = left.exact && right.exact; + match (left.value, right.value) { + (TransformedValue::Predicate(left), TransformedValue::Predicate(right)) => { + TransformedResult::predicate(left.and(right), exact) } - (TransformedResult::Predicate(left), _) => TransformedResult::Predicate(left), - (_, TransformedResult::Predicate(right)) => TransformedResult::Predicate(right), - _ => TransformedResult::NotTransformed, + (TransformedValue::Predicate(left), _) => TransformedResult::predicate(left, false), + (_, TransformedValue::Predicate(right)) => TransformedResult::predicate(right, false), + _ => TransformedResult::not_transformed(), } } fn to_iceberg_or_predicate(left: TransformedResult, right: TransformedResult) -> TransformedResult { - match (left, right) { - (TransformedResult::Predicate(left), TransformedResult::Predicate(right)) => { - TransformedResult::Predicate(left.or(right)) + let exact = left.exact && right.exact; + match (left.value, right.value) { + (TransformedValue::Predicate(left), TransformedValue::Predicate(right)) => { + TransformedResult::predicate(left.or(right), exact) } - _ => TransformedResult::NotTransformed, + _ => TransformedResult::not_transformed(), } } @@ -383,16 +464,29 @@ fn to_iceberg_binary_predicate( right: TransformedResult, op: PredicateOperator, ) -> TransformedResult { - let (r, d, op) = match (left, right) { - (TransformedResult::NotTransformed, _) => return TransformedResult::NotTransformed, - (_, TransformedResult::NotTransformed) => return TransformedResult::NotTransformed, - (TransformedResult::Column(r), TransformedResult::Literal(d)) => (r, d, op), - (TransformedResult::Literal(d), TransformedResult::Column(r)) => { + let exact = left.exact && right.exact; + let (r, d, op) = match (left.value, right.value) { + (TransformedValue::NotTransformed, _) => { + return TransformedResult::not_transformed(); + } + (_, TransformedValue::NotTransformed) => { + return TransformedResult::not_transformed(); + } + (TransformedValue::Column(r), TransformedValue::Literal(d)) => (r, d, op), + (TransformedValue::Literal(d), TransformedValue::Column(r)) => { (r, d, reverse_predicate_operator(op)) } - _ => return TransformedResult::NotTransformed, + _ => return TransformedResult::not_transformed(), }; - TransformedResult::Predicate(Predicate::Binary(BinaryExpression::new(op, r, d))) + // For schema-evolved files that do not contain the referenced column, the + // Arrow reader conservatively matches < and <= predicates. They remain safe + // pruning predicates, but DataFusion must apply the original filter again. + let exact = exact + && !matches!( + op, + PredicateOperator::LessThan | PredicateOperator::LessThanOrEq + ); + TransformedResult::predicate(Predicate::Binary(BinaryExpression::new(op, r, d)), exact) } fn reverse_predicate_operator(op: PredicateOperator) -> PredicateOperator { @@ -442,13 +536,14 @@ mod tests { use datafusion::arrow::datatypes::{DataType, Field, Schema, TimeUnit}; use datafusion::common::DFSchema; + use datafusion::logical_expr::TableProviderFilterPushDown; use datafusion::logical_expr::utils::split_conjunction; use datafusion::prelude::{Expr, SessionContext}; use iceberg::expr::{Predicate, Reference}; use iceberg::spec::Datum; use parquet::arrow::PARQUET_FIELD_ID_META_KEY; - use super::convert_filters_to_predicate; + use super::{classify_filter_pushdown, convert_filters_to_predicate}; fn create_test_schema() -> DFSchema { let arrow_schema = Schema::new(vec![ @@ -480,6 +575,70 @@ mod tests { convert_filters_to_predicate(&exprs[..]) } + fn classify_sql_filter(sql: &str) -> TableProviderFilterPushDown { + let df_schema = create_test_schema(); + let expr = SessionContext::new() + .parse_sql_expr(sql, &df_schema) + .unwrap(); + classify_filter_pushdown(&expr) + } + + #[test] + fn test_filter_pushdown_classification() { + assert_eq!( + classify_sql_filter("foo > 1 AND bar = 'test'"), + TableProviderFilterPushDown::Exact + ); + assert_eq!( + classify_sql_filter("foo > 1 AND length(bar) = 1"), + TableProviderFilterPushDown::Inexact + ); + assert_eq!( + classify_sql_filter("length(bar) = 1"), + TableProviderFilterPushDown::Unsupported + ); + assert_eq!( + classify_sql_filter("foo < 1"), + TableProviderFilterPushDown::Inexact + ); + assert_eq!( + classify_sql_filter("bar NOT IN ('test')"), + TableProviderFilterPushDown::Inexact + ); + assert_eq!( + classify_sql_filter("bar NOT LIKE 'test%'"), + TableProviderFilterPushDown::Inexact + ); + } + + #[test] + fn test_filter_pushdown_with_cast_is_inexact() { + assert_eq!( + classify_sql_filter("ts >= timestamp '2023-01-05T00:00:00'"), + TableProviderFilterPushDown::Inexact + ); + } + + #[test] + fn test_negated_partial_filter_is_unsupported() { + assert_eq!( + classify_sql_filter("NOT (foo > 1 AND length(bar) = 1)"), + TableProviderFilterPushDown::Unsupported + ); + } + + #[test] + fn test_negated_isnan_filter_is_inexact() { + assert_eq!( + classify_sql_filter("isnan(qux)"), + TableProviderFilterPushDown::Exact + ); + assert_eq!( + classify_sql_filter("NOT isnan(qux)"), + TableProviderFilterPushDown::Inexact + ); + } + #[test] fn test_predicate_conversion_with_single_condition() { let predicate = convert_to_iceberg_predicate("foo = 1").unwrap(); diff --git a/crates/integrations/datafusion/src/table/mod.rs b/crates/integrations/datafusion/src/table/mod.rs index 2fd958dff4..63c123d9cd 100644 --- a/crates/integrations/datafusion/src/table/mod.rs +++ b/crates/integrations/datafusion/src/table/mod.rs @@ -50,6 +50,7 @@ use metadata_table::IcebergMetadataTableProvider; use crate::error::to_datafusion_error; use crate::physical_plan::commit::IcebergCommitExec; +use crate::physical_plan::expr_to_predicate::classify_filter_pushdown; use crate::physical_plan::project::project_with_partition; use crate::physical_plan::repartition::repartition; use crate::physical_plan::scan::IcebergTableScan; @@ -146,8 +147,10 @@ impl TableProvider for IcebergTableProvider { &self, filters: &[&Expr], ) -> DFResult> { - // Push down all filters, as a single source of truth, the scanner will drop the filters which couldn't be push down - Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()]) + Ok(filters + .iter() + .map(|filter| classify_filter_pushdown(filter)) + .collect()) } async fn insert_into( @@ -326,8 +329,10 @@ impl TableProvider for IcebergStaticTableProvider { &self, filters: &[&Expr], ) -> DFResult> { - // Push down all filters, as a single source of truth, the scanner will drop the filters which couldn't be push down - Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()]) + Ok(filters + .iter() + .map(|filter| classify_filter_pushdown(filter)) + .collect()) } async fn insert_into( diff --git a/crates/integrations/datafusion/tests/integration_datafusion_test.rs b/crates/integrations/datafusion/tests/integration_datafusion_test.rs index cebac75dd9..214c2df3d5 100644 --- a/crates/integrations/datafusion/tests/integration_datafusion_test.rs +++ b/crates/integrations/datafusion/tests/integration_datafusion_test.rs @@ -274,7 +274,7 @@ async fn test_table_projection() -> Result<()> { } #[tokio::test] -async fn test_table_predict_pushdown() -> Result<()> { +async fn test_table_predicate_pushdown() -> Result<()> { let iceberg_catalog = get_iceberg_catalog().await; let namespace = NamespaceIdent::new("ns".to_string()); set_test_namespace(&iceberg_catalog, &namespace).await?; @@ -315,6 +315,34 @@ async fn test_table_predict_pushdown() -> Result<()> { // the first row is logical_plan, the second row is physical_plan let expected = "predicate:[(foo > 1) OR (bar IS NULL)]"; assert!(s.value(1).trim().contains(expected)); + assert!( + s.value(1).contains("FilterExec"), + "partially converted filters must be re-applied: {}", + s.value(1) + ); + + let records = ctx + .sql("select * from catalog.ns.t1 where foo > 1 and bar is null") + .await + .unwrap() + .explain(false, false) + .unwrap() + .collect() + .await + .unwrap(); + let physical_plan = records[0] + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(1); + + assert!(physical_plan.contains("predicate:[(foo > 1) AND (bar IS NULL)]")); + assert!( + !physical_plan.contains("FilterExec"), + "exact filters should not be re-applied: {physical_plan}" + ); + Ok(()) } diff --git a/crates/sqllogictest/testdata/slts/df_test/binary_predicate_pushdown.slt b/crates/sqllogictest/testdata/slts/df_test/binary_predicate_pushdown.slt index aa68ab2762..611f1c56de 100644 --- a/crates/sqllogictest/testdata/slts/df_test/binary_predicate_pushdown.slt +++ b/crates/sqllogictest/testdata/slts/df_test/binary_predicate_pushdown.slt @@ -22,13 +22,10 @@ query TT EXPLAIN SELECT * FROM default.default.test_binary_table WHERE data = X'0102' ---- -logical_plan -01)Filter: default.default.test_binary_table.data = LargeBinary("1,2") -02)--TableScan: default.default.test_binary_table projection=[id, data], partial_filters=[default.default.test_binary_table.data = LargeBinary("1,2")] +logical_plan TableScan: default.default.test_binary_table projection=[id, data], full_filters=[default.default.test_binary_table.data = LargeBinary("1,2")] physical_plan -01)FilterExec: data@1 = 0102 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,data] predicate:[data = 0102] +01)CooperativeExec +02)--IcebergTableScan projection:[id,data] predicate:[data = 0102] # Verify empty result from empty table query I? diff --git a/crates/sqllogictest/testdata/slts/df_test/boolean_predicate_pushdown.slt b/crates/sqllogictest/testdata/slts/df_test/boolean_predicate_pushdown.slt index 496f719261..a7981e281a 100644 --- a/crates/sqllogictest/testdata/slts/df_test/boolean_predicate_pushdown.slt +++ b/crates/sqllogictest/testdata/slts/df_test/boolean_predicate_pushdown.slt @@ -33,13 +33,10 @@ INSERT INTO default.default.test_boolean_table VALUES query TT EXPLAIN SELECT * FROM default.default.test_boolean_table WHERE is_active = true ---- -logical_plan -01)Filter: default.default.test_boolean_table.is_active -02)--TableScan: default.default.test_boolean_table projection=[id, is_active, description], partial_filters=[default.default.test_boolean_table.is_active] +logical_plan TableScan: default.default.test_boolean_table projection=[id, is_active, description], full_filters=[default.default.test_boolean_table.is_active] physical_plan -01)FilterExec: is_active@1 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,is_active,description] predicate:[is_active = true] +01)CooperativeExec +02)--IcebergTableScan projection:[id,is_active,description] predicate:[is_active = true] # Query with is_active = true query ITT rowsort @@ -53,13 +50,10 @@ SELECT * FROM default.default.test_boolean_table WHERE is_active = true query TT EXPLAIN SELECT * FROM default.default.test_boolean_table WHERE is_active = false ---- -logical_plan -01)Filter: NOT default.default.test_boolean_table.is_active -02)--TableScan: default.default.test_boolean_table projection=[id, is_active, description], partial_filters=[NOT default.default.test_boolean_table.is_active] +logical_plan TableScan: default.default.test_boolean_table projection=[id, is_active, description], full_filters=[NOT default.default.test_boolean_table.is_active] physical_plan -01)FilterExec: NOT is_active@1 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,is_active,description] predicate:[is_active = false] +01)CooperativeExec +02)--IcebergTableScan projection:[id,is_active,description] predicate:[is_active = false] # Query with is_active = false query ITT rowsort @@ -72,13 +66,10 @@ SELECT * FROM default.default.test_boolean_table WHERE is_active = false query TT EXPLAIN SELECT * FROM default.default.test_boolean_table WHERE is_active != true ---- -logical_plan -01)Filter: NOT default.default.test_boolean_table.is_active -02)--TableScan: default.default.test_boolean_table projection=[id, is_active, description], partial_filters=[NOT default.default.test_boolean_table.is_active] +logical_plan TableScan: default.default.test_boolean_table projection=[id, is_active, description], full_filters=[NOT default.default.test_boolean_table.is_active] physical_plan -01)FilterExec: NOT is_active@1 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,is_active,description] predicate:[is_active = false] +01)CooperativeExec +02)--IcebergTableScan projection:[id,is_active,description] predicate:[is_active = false] # Query with is_active != true (includes false and NULL) query ITT rowsort diff --git a/crates/sqllogictest/testdata/slts/df_test/like_predicate_pushdown.slt b/crates/sqllogictest/testdata/slts/df_test/like_predicate_pushdown.slt index 3d8b151aa9..59166fc262 100644 --- a/crates/sqllogictest/testdata/slts/df_test/like_predicate_pushdown.slt +++ b/crates/sqllogictest/testdata/slts/df_test/like_predicate_pushdown.slt @@ -31,13 +31,10 @@ INSERT INTO default.default.test_unpartitioned_table VALUES (5, 'alice'), (6, 'A query TT EXPLAIN SELECT * FROM default.default.test_unpartitioned_table WHERE name LIKE 'Al%' ---- -logical_plan -01)Filter: default.default.test_unpartitioned_table.name LIKE Utf8("Al%") -02)--TableScan: default.default.test_unpartitioned_table projection=[id, name], partial_filters=[default.default.test_unpartitioned_table.name LIKE Utf8("Al%")] +logical_plan TableScan: default.default.test_unpartitioned_table projection=[id, name], full_filters=[default.default.test_unpartitioned_table.name LIKE Utf8("Al%")] physical_plan -01)FilterExec: name@1 LIKE Al% -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,name] predicate:[name STARTS WITH "Al"] +01)CooperativeExec +02)--IcebergTableScan projection:[id,name] predicate:[name STARTS WITH "Al"] # Test LIKE filtering with case-sensitive match query IT rowsort diff --git a/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt b/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt index ffa74173dc..f625bf1ad7 100644 --- a/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt +++ b/crates/sqllogictest/testdata/slts/df_test/timestamp_predicate_pushdown.slt @@ -44,13 +44,10 @@ VALUES query TT EXPLAIN SELECT * FROM default.default.test_timestamp_table WHERE ts = CAST('2023-01-05 12:30:00' AS TIMESTAMP) ---- -logical_plan -01)Filter: default.default.test_timestamp_table.ts = TimestampNanosecond(1672921800000000000, None) -02)--TableScan: default.default.test_timestamp_table projection=[id, ts], partial_filters=[default.default.test_timestamp_table.ts = TimestampNanosecond(1672921800000000000, None)] +logical_plan TableScan: default.default.test_timestamp_table projection=[id, ts], full_filters=[default.default.test_timestamp_table.ts = TimestampNanosecond(1672921800000000000, None)] physical_plan -01)FilterExec: ts@1 = 1672921800000000000 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,ts] predicate:[ts = 2023-01-05 12:30:00] +01)CooperativeExec +02)--IcebergTableScan projection:[id,ts] predicate:[ts = 2023-01-05 12:30:00] # Verify timestamp equality filtering works query I? @@ -62,13 +59,10 @@ SELECT * FROM default.default.test_timestamp_table WHERE ts = CAST('2023-01-05 1 query TT EXPLAIN SELECT * FROM default.default.test_timestamp_table WHERE ts > CAST('2023-01-10 00:00:00' AS TIMESTAMP) ---- -logical_plan -01)Filter: default.default.test_timestamp_table.ts > TimestampNanosecond(1673308800000000000, None) -02)--TableScan: default.default.test_timestamp_table projection=[id, ts], partial_filters=[default.default.test_timestamp_table.ts > TimestampNanosecond(1673308800000000000, None)] +logical_plan TableScan: default.default.test_timestamp_table projection=[id, ts], full_filters=[default.default.test_timestamp_table.ts > TimestampNanosecond(1673308800000000000, None)] physical_plan -01)FilterExec: ts@1 > 1673308800000000000 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,ts] predicate:[ts > 2023-01-10 00:00:00] +01)CooperativeExec +02)--IcebergTableScan projection:[id,ts] predicate:[ts > 2023-01-10 00:00:00] # Verify timestamp greater than filtering query I? rowsort @@ -92,10 +86,10 @@ WHERE ts >= CAST('2023-01-05 00:00:00' AS TIMESTAMP) AND ts <= CAST('2023-01-15 23:59:59' AS TIMESTAMP) ---- logical_plan -01)Filter: default.default.test_timestamp_table.ts >= TimestampNanosecond(1672876800000000000, None) AND default.default.test_timestamp_table.ts <= TimestampNanosecond(1673827199000000000, None) -02)--TableScan: default.default.test_timestamp_table projection=[id, ts], partial_filters=[default.default.test_timestamp_table.ts >= TimestampNanosecond(1672876800000000000, None), default.default.test_timestamp_table.ts <= TimestampNanosecond(1673827199000000000, None)] +01)Filter: default.default.test_timestamp_table.ts <= TimestampNanosecond(1673827199000000000, None) +02)--TableScan: default.default.test_timestamp_table projection=[id, ts], full_filters=[default.default.test_timestamp_table.ts >= TimestampNanosecond(1672876800000000000, None)], partial_filters=[default.default.test_timestamp_table.ts <= TimestampNanosecond(1673827199000000000, None)] physical_plan -01)FilterExec: ts@1 >= 1672876800000000000 AND ts@1 <= 1673827199000000000 +01)FilterExec: ts@1 <= 1673827199000000000 02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 03)----IcebergTableScan projection:[id,ts] predicate:[(ts >= 2023-01-05 00:00:00) AND (ts <= 2023-01-15 23:59:59)] @@ -156,13 +150,10 @@ VALUES query TT EXPLAIN SELECT * FROM default.default.test_timestamp_micros WHERE ts > CAST('2023-01-01 00:00:00' AS TIMESTAMP) ---- -logical_plan -01)Filter: default.default.test_timestamp_micros.ts > TimestampMicrosecond(1672531200000000, None) -02)--TableScan: default.default.test_timestamp_micros projection=[id, ts], partial_filters=[default.default.test_timestamp_micros.ts > TimestampMicrosecond(1672531200000000, None)] +logical_plan TableScan: default.default.test_timestamp_micros projection=[id, ts], full_filters=[default.default.test_timestamp_micros.ts > TimestampMicrosecond(1672531200000000, None)] physical_plan -01)FilterExec: ts@1 > 1672531200000000 -02)--RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 -03)----IcebergTableScan projection:[id,ts] predicate:[ts > 2023-01-01 00:00:00] +01)CooperativeExec +02)--IcebergTableScan projection:[id,ts] predicate:[ts > 2023-01-01 00:00:00] query I? SELECT * FROM default.default.test_timestamp_micros WHERE ts > CAST('2023-01-01 00:00:00' AS TIMESTAMP)