kosiew commented on code in PR #25455:
URL: https://github.com/apache/datafusion/pull/25455#discussion_r4194939366
##########
docs/source/library-user-guide/query-optimizer.md:
##########
@@ -162,6 +162,16 @@ individual operators or expressions.
Sometimes there is an initial pass that visits the plan and builds state that
is used in a second pass that performs
the actual optimization. This approach is used in projection push down and
filter push down.
+### Rule Precedence
+
+Two rules can want the opposite order for the same pair of adjacent plan
nodes. Each rule then undoes the work of the other one on every optimizer pass.
The loop stops at `datafusion.optimizer.max_passes`, and the position of the
two rules in the rule list decides the plan. This wastes a pass and makes the
plan depend on the rule list.
+
+Do not let two rules compete. Give one rule precedence, make that rule yield,
and record the decision in the module documentation of both rules.
Review Comment:
I think this should say "Give one rule precedence; make the competing rule
yield," since the rule with precedence is the one that wins. Also, the
optimizer can stop early when an end-of-pass plan signature repeats, so please
avoid saying that the loop necessarily runs until `max_passes`.
##########
datafusion/sqllogictest/test_files/projection_pushdown.slt:
##########
@@ -2331,12 +2331,232 @@ logical_plan
05)--------SubqueryAlias: inner_t
06)----------Projection: simple_struct.id, simple_struct.s,
__datafusion_extracted_1
07)------------Limit: skip=0, fetch=1
-08)--------------Filter: __datafusion_extracted_2 > Int64(120) AND
__datafusion_extracted_1 != Utf8("delta")
-09)----------------Filter: simple_struct.id = outer_ref(outer_t.id)
-10)------------------Projection: get_field(simple_struct.s, Utf8("value")) AS
__datafusion_extracted_2, get_field(simple_struct.s, Utf8("label")) AS
__datafusion_extracted_1, simple_struct.id, simple_struct.s
-11)--------------------TableScan: simple_struct,
partial_filters=[simple_struct.id = outer_ref(outer_t.id),
get_field(simple_struct.s, Utf8("value")) > Int64(120)]
-12)--SubqueryAlias: outer_t
-13)----TableScan: simple_struct projection=[id]
+08)--------------Filter: simple_struct.id = outer_ref(outer_t.id) AND
__datafusion_extracted_2 > Int64(120) AND __datafusion_extracted_1 !=
Utf8("delta")
+09)----------------Projection: get_field(simple_struct.s, Utf8("value")) AS
__datafusion_extracted_2, simple_struct.id, simple_struct.s,
get_field(simple_struct.s, Utf8("label")) AS __datafusion_extracted_1
+10)------------------TableScan: simple_struct,
partial_filters=[simple_struct.id = outer_ref(outer_t.id),
get_field(simple_struct.s, Utf8("value")) > Int64(120)]
+11)--SubqueryAlias: outer_t
+12)----TableScan: simple_struct projection=[id]
statement ok
set datafusion.explain.logical_plan_only = false;
+
+#####################
+# Section 17: Precedence between PushDownFilter and PushDownLeafProjections
+#
+# Both rules move nodes towards the leaves, and for an adjacent filter and
+# "pure extraction projection" they want the opposite order. The precedence is:
+# the extraction projection wins. It stays next to the scan, so the scan can
+# absorb it and read only the struct leaf, and the filter stays above it.
+#
+# See https://github.com/apache/datafusion/issues/14540
+#####################
+
+statement ok
+SET datafusion.execution.target_partitions = 1;
+
+statement ok
+CREATE TABLE events_mem (
+ "date" DATE,
+ "timestamp" TIMESTAMP,
+ ids STRUCT<id1 VARCHAR, extra INT>,
+ structs STRUCT<var1 VARCHAR, extra VARCHAR>
+) AS VALUES
+ (DATE '2025-01-03', TIMESTAMP '2025-01-03 01:00:00', {id1: 'dev1', extra:
1}, {var1: 'user1', extra: 'e1'}),
+ (DATE '2025-01-03', TIMESTAMP '2025-01-03 02:00:00', {id1: 'dev1', extra:
2}, {var1: 'user2', extra: 'e2'}),
+ (DATE '2025-01-04', TIMESTAMP '2025-01-04 01:00:00', {id1: 'dev2', extra:
3}, {var1: 'user3', extra: 'e3'});
+
+# The query from #14540. The extraction projection is the bottom node and the
+# filter sits above it.
+query TT
+EXPLAIN WITH events AS (
+ SELECT ids.id1 AS device, structs.var1 AS user, "timestamp"
+ FROM events_mem
+ WHERE "date" = '2025-01-03'
+)
+SELECT *, LAG(user, 1) OVER (PARTITION BY device ORDER BY "timestamp") AS prev
+FROM events
+WHERE device IS NOT NULL AND device != '' AND user IS NOT NULL AND user != ''
+LIMIT 100;
+----
+logical_plan
+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
+02)--Limit: skip=0, fetch=100
+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]]
+04)------SubqueryAlias: events
+05)--------Projection: __datafusion_extracted_1 AS device,
__datafusion_extracted_2 AS user, events_mem.timestamp
+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("")
+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
+08)--------------TableScan: events_mem projection=[date, timestamp, ids,
structs]
+physical_plan
+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]
+02)--GlobalLimitExec: skip=0, fetch=100
+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]
+04)------SortExec: expr=[device@0 ASC NULLS LAST, timestamp@2 ASC NULLS LAST],
preserve_partitioning=[false]
+05)--------ProjectionExec: expr=[__datafusion_extracted_1@0 as device,
__datafusion_extracted_2@1 as user, timestamp@2 as timestamp]
+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]
+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]
+08)--------------DataSourceExec: partitions=1, partition_sizes=[1]
+
+# Determinism: the plan is a fixed point reached before the default pass limit,
Review Comment:
These pass-limit checks show that the final plans match, but they do not
prove that individual rules have stopped undoing each other within a pass. As
optional extra coverage, could you add an observer-based test that verifies
each rule leaves the plan unchanged during the final pass?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]