Skip to content

Filter pushdown picks the wrong column when a join output has duplicated column names, causing wrong results with TopK dynamic filters #25296

Description

@jayzhan211

Describe the bug

Physical filter pushdown resolves columns by name, so when a schema contains two columns with the same name, a filter over the second one is pushed to the first. With a TopK (ORDER BY ... LIMIT n) dynamic filter this returns wrong results under the default configuration.

Trigger: two joined tables that share a column name. Below, orders and payments both have amount, so the lower join's output schema is [amount, id, amount]. A filter over payments.amount is looked up as "amount" in that schema, matches the first amount, which is orders.amount, and is pushed to the wrong table. If either column were renamed, the query would be correct.

Query: the order with the smallest payment.

SELECT o.amount, p.amount
FROM orders o
JOIN payments p ON o.id = p.id
JOIN customers c ON p.id = c.id
ORDER BY p.amount LIMIT 1;

How the filter travels down the physical plan (the lower join is built on payments because orders is larger, and a projection restores the written column order):

SortExec TopK(fetch=1), ORDER BY amount@1          dynamic filter: amount@1 < 30
   │                                               ("amount" at output index 1 = p.amount, correct)
   ▼
HashJoinExec (upper)  output [o.amount, p.amount]
   │   index 1 is on the right child, so the filter may go there.
   │   Position inside the child is resolved BY NAME: index_of("amount")
   ▼
HashJoinExec (lower)  output [amount,  id,  amount]      <-- two columns named "amount"
                               ▲             ▲
                               │             └── p.amount   the column the filter is about
                               └──────────────── o.amount   the column index_of("amount") returns:
                                                            same name, first match wins
   │
   ├── DataSourceExec payments      receives nothing        (should have received amount < 30)
   └── DataSourceExec orders        receives amount < 30    WRONG TABLE

Why that returns the wrong row. The joined rows are (100, 30), (200, 20), (300, 10), and the answer should be (300, 10). orders streams through the lower join one row per batch (batch_size = 1, one row per row group), and after each batch the TopK publishes "only rows with amount below my current best can still matter":

step orders scan row reaching the TopK TopK best so far filter published by the TopK
1 reads order 1 (amount 100) (100, 30) 30 amount < 30
2 order 2 (amount 200): 200 < 30 is false, row group pruned nothing 30 unchanged
3 order 3 (amount 300): 300 < 30 is false, pruned nothing 30 unchanged
4-6 orders 4 to 6 pruned the same way nothing 30 unchanged

The threshold 30 was computed from payments.amount (values 30, 20, 10). On that column amount < 30 keeps payments 20 and 10, the ones that could still beat the current best. But it is applied to orders.amount (values 100, 200, ..., 600), where no value is below 30. The parquet reader checks each remaining row group's min/max against the filter before decoding it, finds 200 < 30, 300 < 30, ... impossible, and skips them all. The rows (200, 20) and (300, 10) never reach the TopK, which reports the only row it ever saw:

                    expected     actual
o.amount, p.amount  300, 10      100, 30

Root cause: FilterRemapper::try_remap in datafusion/physical-plan/src/filter_pushdown.rs maps a parent filter into a child's schema with child_schema.index_of(col.name()). Arrow's index_of returns the first field with that name, so with duplicated names the lookup is ambiguous and silently picks the wrong column. HashJoinExec::gather_filters_for_pushdown reaches it through ChildFilterDescription::from_child_with_allowed_indices; the index check there only decides whether a column belongs to the child, the position inside the child is still looked up by name. Physical schemas are not required to have unique names, and a join of two tables that share a column name produces one every time.

To Reproduce

set datafusion.execution.target_partitions = 1;
-- one row per batch, so the TopK filter is updated while the `orders` scan still has row groups to read
set datafusion.execution.batch_size = 1;

CREATE TABLE orders_src(id INT, amount INT) AS VALUES (1, 100), (2, 200), (3, 300), (4, 400), (5, 500), (6, 600);
CREATE TABLE payments_src(id INT, amount INT) AS VALUES (1, 30), (2, 20), (3, 10);
CREATE TABLE customers_src(id INT) AS VALUES (1), (2), (3);

COPY orders_src TO '/tmp/dup/orders.parquet' STORED AS PARQUET OPTIONS ('format.max_row_group_size' '1');
COPY payments_src TO '/tmp/dup/payments.parquet' STORED AS PARQUET;
COPY customers_src TO '/tmp/dup/customers.parquet' STORED AS PARQUET;

CREATE EXTERNAL TABLE orders(id INT, amount INT) STORED AS PARQUET LOCATION '/tmp/dup/orders.parquet';
CREATE EXTERNAL TABLE payments(id INT, amount INT) STORED AS PARQUET LOCATION '/tmp/dup/payments.parquet';
CREATE EXTERNAL TABLE customers(id INT) STORED AS PARQUET LOCATION '/tmp/dup/customers.parquet';

SELECT o.amount, p.amount
FROM orders o
JOIN payments p ON o.id = p.id
JOIN customers c ON p.id = c.id
ORDER BY p.amount LIMIT 1;
+--------+--------+
| amount | amount |
+--------+--------+
| 100    | 30     |
+--------+--------+

The same wrong row comes back with datafusion.execution.parquet.pushdown_filters = true. With datafusion.optimizer.enable_topk_dynamic_filter_pushdown = false the result is correct.

Expected behavior

+--------+--------+
| amount | amount |
+--------+--------+
| 300    | 10     |
+--------+--------+

Additional context

EXPLAIN of the query:

SortExec: TopK(fetch=1), expr=[amount@1 ASC NULLS LAST], preserve_partitioning=[false]
--HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(id@0, id@1)], projection=[amount@1, amount@3]
----DataSourceExec: file_groups={1 group: [[.../customers.parquet]]}, projection=[id], file_type=parquet
----HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(id@0, id@0)], projection=[amount@3, id@0, amount@1]
------DataSourceExec: file_groups={1 group: [[.../payments.parquet]]}, projection=[id, amount], file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible
------DataSourceExec: file_groups={1 group: [[.../orders.parquet]]}, projection=[id, amount], file_type=parquet, predicate=DynamicFilter [ empty ] AND DynamicFilter [ empty ], dynamic_rg_pruning=eligible

EXPLAIN ANALYZE shows the misplaced filter on the orders scan: predicate=DynamicFilter [ ... ] AND DynamicFilter [ amount@1 < 30 ].

Why the reproducer needs this shape: the name lookup only misfires when a join's child is itself a join (two same-named columns in one schema), and the wrongly targeted table must be the streaming probe side, so that the filter still has row groups left to prune when the TopK updates it. With the default batch size the lower join coalesces its output into one batch and the bug is masked; batch_size = 1 and one row per row group make it deterministic.

Found while reviewing #25255, which removes one caller of the name-based semi-join routing but does not change try_remap. Reproduced on main at a407990 and on that PR branch. The static-predicate path is less exposed only because the logical optimizer usually pushes static predicates to the scans before the physical pass.

A fix probably needs the remap to carry positions rather than names: HashJoinExec already knows which parent output index maps to which child index (column_indices), so from_child_with_allowed_indices could take a HashMap<usize, usize> from parent index to child index instead of a HashSet<usize> plus a name lookup.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething isn't working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions