Skip to content

Commit f146f1d

Browse files
adriangbclaude
andcommitted
test: cover the filter and extraction projection precedence
Unit tests in `push_down_filter.rs`: a filter stays above a pure extraction projection, and a projection that also computes an expression is still pushed through. A new section in `projection_pushdown.slt` with the query from #14540 on a memory table and on a Parquet file. The Parquet plan shows both properties the precedence must keep: `DataSourceExec` reads only the struct leaf `ids.id1`, and the `date` predicate reaches the scan for row group pruning. The section also pins the plan at a reduced `datafusion.optimizer.max_passes`. The simple shape gives the same plan at one pass as at the default. The issue shape gives the same plan at two passes, because the two filters merge in the second pass. Neither plan depends on the pass limit, so neither depends on which of the two rules runs last. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
1 parent 81c4d05 commit f146f1d

2 files changed

Lines changed: 285 additions & 0 deletions

File tree

‎datafusion/optimizer/src/push_down_filter.rs‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4590,6 +4590,70 @@ mod tests {
45904590
)
45914591
}
45924592

4593+
/// A filter is not moved below a pure extraction projection, even when its
4594+
/// predicate only references pass-through columns.
4595+
///
4596+
/// `PushDownLeafProjections` moves such a projection back below the filter,
4597+
/// so pushing here would make the two rules undo each other on every
4598+
/// optimizer pass. See <https://github.com/apache/datafusion/issues/14540>
4599+
/// and the module documentation.
4600+
#[test]
4601+
fn filter_not_pushed_through_pure_extraction_projection() -> Result<()> {
4602+
let table_scan = test_table_scan()?;
4603+
4604+
// The shape `ExtractLeafExpressions` produces: one extraction alias
4605+
// plus pass-through columns.
4606+
let proj = LogicalPlanBuilder::from(table_scan)
4607+
.project(vec![
4608+
leaf_udf_expr(col("a")).alias("__datafusion_extracted_1"),
4609+
col("b"),
4610+
col("c"),
4611+
])?
4612+
.build()?;
4613+
4614+
// `b` is a plain pass-through column, so without the precedence rule
4615+
// this predicate would reach the scan as a `full_filters` entry.
4616+
let plan = LogicalPlanBuilder::from(proj)
4617+
.filter(col("b").gt(lit(5i64)))?
4618+
.build()?;
4619+
4620+
assert_optimized_plan_equal!(
4621+
plan,
4622+
@r"
4623+
Filter: test.b > Int64(5)
4624+
Projection: leaf_udf(test.a) AS __datafusion_extracted_1, test.b, test.c
4625+
TableScan: test
4626+
"
4627+
)
4628+
}
4629+
4630+
/// A projection that mixes an extraction alias with a computed expression is
4631+
/// not a pure extraction projection, so the filter still moves below it.
4632+
#[test]
4633+
fn filter_pushed_through_mixed_extraction_projection() -> Result<()> {
4634+
let table_scan = test_table_scan()?;
4635+
4636+
let proj = LogicalPlanBuilder::from(table_scan)
4637+
.project(vec![
4638+
leaf_udf_expr(col("a")).alias("__datafusion_extracted_1"),
4639+
(col("b") + lit(1i64)).alias("b_plus"),
4640+
col("c"),
4641+
])?
4642+
.build()?;
4643+
4644+
let plan = LogicalPlanBuilder::from(proj)
4645+
.filter(col("c").gt(lit(5i64)))?
4646+
.build()?;
4647+
4648+
assert_optimized_plan_equal!(
4649+
plan,
4650+
@r"
4651+
Projection: leaf_udf(test.a) AS __datafusion_extracted_1, test.b + Int64(1) AS b_plus, test.c
4652+
TableScan: test, full_filters=[test.c > Int64(5)]
4653+
"
4654+
)
4655+
}
4656+
45934657
#[test]
45944658
fn filter_not_pushed_down_through_table_scan_with_fetch() -> Result<()> {
45954659
let scan = test_table_scan()?;

‎datafusion/sqllogictest/test_files/projection_pushdown.slt‎

Lines changed: 221 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2339,3 +2339,224 @@ logical_plan
23392339

23402340
statement ok
23412341
set datafusion.explain.logical_plan_only = false;
2342+
2343+
#####################
2344+
# Section 17: Precedence between PushDownFilter and PushDownLeafProjections
2345+
#
2346+
# Both rules move nodes towards the leaves, and for an adjacent filter and
2347+
# "pure extraction projection" they want the opposite order. The precedence is:
2348+
# the extraction projection wins. It stays next to the scan, so the scan can
2349+
# absorb it and read only the struct leaf, and the filter stays above it.
2350+
#
2351+
# See https://github.com/apache/datafusion/issues/14540
2352+
#####################
2353+
2354+
statement ok
2355+
SET datafusion.execution.target_partitions = 1;
2356+
2357+
statement ok
2358+
CREATE TABLE events_mem (
2359+
"date" DATE,
2360+
"timestamp" TIMESTAMP,
2361+
ids STRUCT<id1 VARCHAR, extra INT>,
2362+
structs STRUCT<var1 VARCHAR, extra VARCHAR>
2363+
) AS VALUES
2364+
(DATE '2025-01-03', TIMESTAMP '2025-01-03 01:00:00', {id1: 'dev1', extra: 1}, {var1: 'user1', extra: 'e1'}),
2365+
(DATE '2025-01-03', TIMESTAMP '2025-01-03 02:00:00', {id1: 'dev1', extra: 2}, {var1: 'user2', extra: 'e2'}),
2366+
(DATE '2025-01-04', TIMESTAMP '2025-01-04 01:00:00', {id1: 'dev2', extra: 3}, {var1: 'user3', extra: 'e3'});
2367+
2368+
# The query from #14540. The extraction projection is the bottom node and the
2369+
# filter sits above it.
2370+
query TT
2371+
EXPLAIN WITH events AS (
2372+
SELECT ids.id1 AS device, structs.var1 AS user, "timestamp"
2373+
FROM events_mem
2374+
WHERE "date" = '2025-01-03'
2375+
)
2376+
SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev
2377+
FROM events
2378+
WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != ''
2379+
LIMIT 100;
2380+
----
2381+
logical_plan
2382+
01)Projection: events.device, events.user, events.timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS prev
2383+
02)--Limit: skip=0, fetch=100
2384+
03)----WindowAggr: windowExpr=[[lag(events.user, Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
2385+
04)------SubqueryAlias: events
2386+
05)--------Projection: __datafusion_extracted_1 AS device, __datafusion_extracted_2 AS user, events_mem.timestamp
2387+
06)----------Filter: events_mem.date = Date32("2025-01-03") AND __datafusion_extracted_1 IS NOT NULL AND __datafusion_extracted_1 != Utf8View("") AND __datafusion_extracted_2 IS NOT NULL AND __datafusion_extracted_2 != Utf8View("")
2388+
07)------------Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, get_field(events_mem.structs, Utf8("var1")) AS __datafusion_extracted_2, events_mem.date, events_mem.timestamp
2389+
08)--------------TableScan: events_mem projection=[date, timestamp, ids, structs]
2390+
physical_plan
2391+
01)ProjectionExec: expr=[device@0 as device, user@1 as user, timestamp@2 as timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as prev]
2392+
02)--GlobalLimitExec: skip=0, fetch=100
2393+
03)----BoundedWindowAggExec: wdw=[lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Utf8View }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
2394+
04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST], preserve_partitioning=[false]
2395+
05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device, __datafusion_extracted_2@1 as user, timestamp@2 as timestamp]
2396+
06)----------FilterExec: date@2 = 2025-01-03 AND __datafusion_extracted_1@0 IS NOT NULL AND __datafusion_extracted_1@0 != AND __datafusion_extracted_2@1 IS NOT NULL AND __datafusion_extracted_2@1 != , projection=[__datafusion_extracted_1@0, __datafusion_extracted_2@1, timestamp@3]
2397+
07)------------ProjectionExec: expr=[get_field(ids@2, id1) as __datafusion_extracted_1, get_field(structs@3, var1) as __datafusion_extracted_2, date@0 as date, timestamp@1 as timestamp]
2398+
08)--------------DataSourceExec: partitions=1, partition_sizes=[1]
2399+
2400+
# Determinism: the plan is a fixed point reached before the default pass limit,
2401+
# so it does not depend on which of the two rules runs last. One pass is enough
2402+
# for the simple shape.
2403+
statement ok
2404+
SET datafusion.optimizer.max_passes = 1;
2405+
2406+
query TT
2407+
EXPLAIN SELECT ids['id1'] FROM events_mem WHERE "date" = '2025-01-03';
2408+
----
2409+
logical_plan
2410+
01)Projection: __datafusion_extracted_1 AS events_mem.ids[id1]
2411+
02)--Filter: events_mem.date = Date32("2025-01-03")
2412+
03)----Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, events_mem.date
2413+
04)------TableScan: events_mem projection=[date, ids]
2414+
physical_plan
2415+
01)ProjectionExec: expr=[__datafusion_extracted_1@0 as events_mem.ids[id1]]
2416+
02)--FilterExec: date@1 = 2025-01-03, projection=[__datafusion_extracted_1@0]
2417+
03)----ProjectionExec: expr=[get_field(ids@1, id1) as __datafusion_extracted_1, date@0 as date]
2418+
04)------DataSourceExec: partitions=1, partition_sizes=[1]
2419+
2420+
statement ok
2421+
RESET datafusion.optimizer.max_passes;
2422+
2423+
query TT
2424+
EXPLAIN SELECT ids['id1'] FROM events_mem WHERE "date" = '2025-01-03';
2425+
----
2426+
logical_plan
2427+
01)Projection: __datafusion_extracted_1 AS events_mem.ids[id1]
2428+
02)--Filter: events_mem.date = Date32("2025-01-03")
2429+
03)----Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, events_mem.date
2430+
04)------TableScan: events_mem projection=[date, ids]
2431+
physical_plan
2432+
01)ProjectionExec: expr=[__datafusion_extracted_1@0 as events_mem.ids[id1]]
2433+
02)--FilterExec: date@1 = 2025-01-03, projection=[__datafusion_extracted_1@0]
2434+
03)----ProjectionExec: expr=[get_field(ids@1, id1) as __datafusion_extracted_1, date@0 as date]
2435+
04)------DataSourceExec: partitions=1, partition_sizes=[1]
2436+
2437+
# The #14540 shape needs a second pass, because the two filters merge into one.
2438+
# It is stable from that pass on.
2439+
statement ok
2440+
SET datafusion.optimizer.max_passes = 2;
2441+
2442+
query TT
2443+
EXPLAIN WITH events AS (
2444+
SELECT ids.id1 AS device, structs.var1 AS user, "timestamp"
2445+
FROM events_mem
2446+
WHERE "date" = '2025-01-03'
2447+
)
2448+
SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev
2449+
FROM events
2450+
WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != ''
2451+
LIMIT 100;
2452+
----
2453+
logical_plan
2454+
01)Projection: events.device, events.user, events.timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS prev
2455+
02)--Limit: skip=0, fetch=100
2456+
03)----WindowAggr: windowExpr=[[lag(events.user, Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
2457+
04)------SubqueryAlias: events
2458+
05)--------Projection: __datafusion_extracted_1 AS device, __datafusion_extracted_2 AS user, events_mem.timestamp
2459+
06)----------Filter: events_mem.date = Date32("2025-01-03") AND __datafusion_extracted_1 IS NOT NULL AND __datafusion_extracted_1 != Utf8View("") AND __datafusion_extracted_2 IS NOT NULL AND __datafusion_extracted_2 != Utf8View("")
2460+
07)------------Projection: get_field(events_mem.ids, Utf8("id1")) AS __datafusion_extracted_1, get_field(events_mem.structs, Utf8("var1")) AS __datafusion_extracted_2, events_mem.date, events_mem.timestamp
2461+
08)--------------TableScan: events_mem projection=[date, timestamp, ids, structs]
2462+
physical_plan
2463+
01)ProjectionExec: expr=[device@0 as device, user@1 as user, timestamp@2 as timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as prev]
2464+
02)--GlobalLimitExec: skip=0, fetch=100
2465+
03)----BoundedWindowAggExec: wdw=[lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Utf8View }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
2466+
04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST], preserve_partitioning=[false]
2467+
05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device, __datafusion_extracted_2@1 as user, timestamp@2 as timestamp]
2468+
06)----------FilterExec: date@2 = 2025-01-03 AND __datafusion_extracted_1@0 IS NOT NULL AND __datafusion_extracted_1@0 != AND __datafusion_extracted_2@1 IS NOT NULL AND __datafusion_extracted_2@1 != , projection=[__datafusion_extracted_1@0, __datafusion_extracted_2@1, timestamp@3]
2469+
07)------------ProjectionExec: expr=[get_field(ids@2, id1) as __datafusion_extracted_1, get_field(structs@3, var1) as __datafusion_extracted_2, date@0 as date, timestamp@1 as timestamp]
2470+
08)--------------DataSourceExec: partitions=1, partition_sizes=[1]
2471+
2472+
statement ok
2473+
RESET datafusion.optimizer.max_passes;
2474+
2475+
query TTT
2476+
WITH events AS (
2477+
SELECT ids.id1 AS device, structs.var1 AS user, "timestamp"
2478+
FROM events_mem
2479+
WHERE "date" = '2025-01-03'
2480+
)
2481+
SELECT device, user, CAST(LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS VARCHAR) AS prev
2482+
FROM events
2483+
WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != ''
2484+
ORDER BY user;
2485+
----
2486+
dev1 user1 NULL
2487+
dev1 user2 user1
2488+
2489+
statement ok
2490+
DROP TABLE events_mem;
2491+
2492+
# Parquet: both properties must survive. `DataSourceExec` reads only the struct
2493+
# leaf `ids.id1` (the `get_field` is in the scan projection), and the `date`
2494+
# predicate reaches the scan for row group pruning.
2495+
statement ok
2496+
COPY (
2497+
SELECT
2498+
column1 AS "date",
2499+
column2 AS "timestamp",
2500+
column3 AS ids,
2501+
column4 AS structs
2502+
FROM VALUES
2503+
(DATE '2025-01-03', TIMESTAMP '2025-01-03 01:00:00', {id1: 'dev1', extra: 1}, {var1: 'user1', extra: 'e1'}),
2504+
(DATE '2025-01-03', TIMESTAMP '2025-01-03 02:00:00', {id1: 'dev1', extra: 2}, {var1: 'user2', extra: 'e2'}),
2505+
(DATE '2025-01-04', TIMESTAMP '2025-01-04 01:00:00', {id1: 'dev2', extra: 3}, {var1: 'user3', extra: 'e3'})
2506+
) TO 'test_files/scratch/projection_pushdown/events.parquet'
2507+
STORED AS PARQUET;
2508+
2509+
statement ok
2510+
CREATE EXTERNAL TABLE events_pq STORED AS PARQUET
2511+
LOCATION 'test_files/scratch/projection_pushdown/events.parquet';
2512+
2513+
query TT
2514+
EXPLAIN SELECT ids['id1'] FROM events_pq WHERE "date" = '2025-01-03';
2515+
----
2516+
logical_plan
2517+
01)Projection: __datafusion_extracted_1 AS events_pq.ids[id1]
2518+
02)--Filter: events_pq.date = Date32("2025-01-03")
2519+
03)----Projection: get_field(events_pq.ids, Utf8("id1")) AS __datafusion_extracted_1, events_pq.date
2520+
04)------TableScan: events_pq projection=[date, ids], partial_filters=[events_pq.date = Date32("2025-01-03")]
2521+
physical_plan
2522+
01)ProjectionExec: expr=[__datafusion_extracted_1@0 as events_pq.ids[id1]]
2523+
02)--FilterExec: date@1 = 2025-01-03, projection=[__datafusion_extracted_1@0]
2524+
03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/projection_pushdown/events.parquet]]}, projection=[get_field(ids@2, id1) as __datafusion_extracted_1, date], file_type=parquet, predicate=date@0 = 2025-01-03, pruning_predicate=date_null_count@2 != row_count@3 AND date_min@0 <= 2025-01-03 AND 2025-01-03 <= date_max@1, required_guarantees=[date in (2025-01-03)]
2525+
2526+
query TT
2527+
EXPLAIN WITH events AS (
2528+
SELECT ids.id1 AS device, structs.var1 AS user, "timestamp"
2529+
FROM events_pq
2530+
WHERE "date" = '2025-01-03'
2531+
)
2532+
SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev
2533+
FROM events
2534+
WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != ''
2535+
LIMIT 100;
2536+
----
2537+
logical_plan
2538+
01)Projection: events.device, events.user, events.timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS prev
2539+
02)--Limit: skip=0, fetch=100
2540+
03)----WindowAggr: windowExpr=[[lag(events.user, Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
2541+
04)------SubqueryAlias: events
2542+
05)--------Projection: __datafusion_extracted_1 AS device, __datafusion_extracted_2 AS user, events_pq.timestamp
2543+
06)----------Filter: events_pq.date = Date32("2025-01-03") AND __datafusion_extracted_1 IS NOT NULL AND __datafusion_extracted_1 != Utf8("") AND __datafusion_extracted_2 IS NOT NULL AND __datafusion_extracted_2 != Utf8("")
2544+
07)------------Projection: get_field(events_pq.ids, Utf8("id1")) AS __datafusion_extracted_1, get_field(events_pq.structs, Utf8("var1")) AS __datafusion_extracted_2, events_pq.date, events_pq.timestamp
2545+
08)--------------TableScan: events_pq projection=[date, timestamp, ids, structs], partial_filters=[events_pq.date = Date32("2025-01-03")]
2546+
physical_plan
2547+
01)ProjectionExec: expr=[device@0 as device, user@1 as user, timestamp@2 as timestamp, lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as prev]
2548+
02)--GlobalLimitExec: skip=0, fetch=100
2549+
03)----BoundedWindowAggExec: wdw=[lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "lag(events.user,Int64(1)) PARTITION BY [events.device] ORDER BY [events.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": nullable Utf8 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
2550+
04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST], preserve_partitioning=[false]
2551+
05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device, __datafusion_extracted_2@1 as user, timestamp@2 as timestamp]
2552+
06)----------FilterExec: date@2 = 2025-01-03 AND __datafusion_extracted_1@0 IS NOT NULL AND __datafusion_extracted_1@0 != AND __datafusion_extracted_2@1 IS NOT NULL AND __datafusion_extracted_2@1 != , projection=[__datafusion_extracted_1@0, __datafusion_extracted_2@1, timestamp@3]
2553+
07)------------DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/projection_pushdown/events.parquet]]}, projection=[get_field(ids@2, id1) as __datafusion_extracted_1, get_field(structs@3, var1) as __datafusion_extracted_2, date, timestamp], file_type=parquet, predicate=date@0 = 2025-01-03 AND get_field(ids@2, id1) IS NOT NULL AND get_field(ids@2, id1) != AND get_field(structs@3, var1) IS NOT NULL AND get_field(structs@3, var1) != , pruning_predicate=date_null_count@2 != row_count@3 AND date_min@0 <= 2025-01-03 AND 2025-01-03 <= date_max@1, required_guarantees=[date in (2025-01-03)]
2554+
2555+
statement ok
2556+
DROP TABLE events_pq;
2557+
2558+
statement ok
2559+
RESET datafusion.optimizer.max_passes;
2560+
2561+
statement ok
2562+
SET datafusion.execution.target_partitions = 4;

0 commit comments

Comments
 (0)