Skip to content

fix(lambda): only push referenced params into the merged batch (#24162) - #166

Merged
LiaCastaneda merged 2 commits into
branch-54from
lia/cherry-pick-lambda-fix
Aug 12, 2026
Merged

fix(lambda): only push referenced params into the merged batch (#24162)#166
LiaCastaneda merged 2 commits into
branch-54from
lia/cherry-pick-lambda-fix

Conversation

@LiaCastaneda

Copy link
Copy Markdown

Cherry picks apache#24162

…e#24162)

## Which issue does this PR close?

basically this PR apache#22853 + a
few more tests

## Rationale for this change

The current lambdas in DF only take a single parameter `(v -> ...)`, so
nobody had noticed that `LambdaExpr` mishandles lambdas with more than
one parameter. The bug surfaced while working on `transform_values`
(apache#22689), which needs `(k, v) -> expr ` two parameters, one of which is
very often unused (e.g. `(k, v) -> v * 2`, k never referenced).

The bug is that when a higher order function with more than 1 param
evaluates a lambda, it fills each parameter into a slot based on its
declared position — for example for `(k, v) -> v` `k` always goes into
slot 0, `v` always into slot 1. `LambdaExpr` separately scans the body
and renumbers whatever it finds referenced into a dense `0..n` range, to
avoid carrying around columns nothing uses (like `v` in this case). That
renumbering is fine for outer captures, but applying it to the lambda's
own parameters is wrong, because it changes where the body looks for a
value without changing where the evaluator put it.

### Example:

in `(k, v) -> v` `v` is declared second (slot 1), but since it's the
only parameter the body references, the renumbering logic reassigns it
to slot 0. The evaluator, unaware of this, writes `k`'s values into slot
0 and `v`'s into slot 1. So the body ends up reading slot 0 expecting
`v` — and gets `k` instead. So the results end up being incorrect.

## What changes are included in this PR?

- `LambdaExpr` now computes `used_params`: which is the subset of its
own declared parameters that are actually referenced in the body.
- `LambdaArgument::new` takes `used_params` and only pushes the
referenced parameters in the body into the merged batch, in original
declaration order — so the body's indices always line up with what's
actually built.
- `HigherOrderFunctionExpr::evaluate` forwards `lambda.used_params()` to
`LambdaArgument::new`

## Are these changes tested?

yes, added two new tests one for the unused-parameter case and
nested-lambda for the shadowing case.

## Are there any user-facing changes?

The only public api change is on `LambdaArgument::new ` which now
requires a new argument: `used_params: &HashSet<String>`, however
LambdaArgument::new is very unlikely to be called outside datafusion,
see
[this](apache#22853 (comment))
comment

(cherry picked from commit 4e6acfe)
@LiaCastaneda
LiaCastaneda merged commit 79de5e9 into branch-54 Aug 12, 2026
67 checks passed
@LiaCastaneda
LiaCastaneda deleted the lia/cherry-pick-lambda-fix branch August 12, 2026 08:54
jayshrivastava pushed a commit that referenced this pull request Aug 13, 2026
…e#24162) (#166)

* fix(lambda): only push referenced params into the merged batch (apache#24162)

## Which issue does this PR close?

basically this PR apache#22853 + a
few more tests

## Rationale for this change

The current lambdas in DF only take a single parameter `(v -> ...)`, so
nobody had noticed that `LambdaExpr` mishandles lambdas with more than
one parameter. The bug surfaced while working on `transform_values`
(apache#22689), which needs `(k, v) -> expr ` two parameters, one of which is
very often unused (e.g. `(k, v) -> v * 2`, k never referenced).

The bug is that when a higher order function with more than 1 param
evaluates a lambda, it fills each parameter into a slot based on its
declared position — for example for `(k, v) -> v` `k` always goes into
slot 0, `v` always into slot 1. `LambdaExpr` separately scans the body
and renumbers whatever it finds referenced into a dense `0..n` range, to
avoid carrying around columns nothing uses (like `v` in this case). That
renumbering is fine for outer captures, but applying it to the lambda's
own parameters is wrong, because it changes where the body looks for a
value without changing where the evaluator put it.

### Example:

in `(k, v) -> v` `v` is declared second (slot 1), but since it's the
only parameter the body references, the renumbering logic reassigns it
to slot 0. The evaluator, unaware of this, writes `k`'s values into slot
0 and `v`'s into slot 1. So the body ends up reading slot 0 expecting
`v` — and gets `k` instead. So the results end up being incorrect.

## What changes are included in this PR?

- `LambdaExpr` now computes `used_params`: which is the subset of its
own declared parameters that are actually referenced in the body.
- `LambdaArgument::new` takes `used_params` and only pushes the
referenced parameters in the body into the merged batch, in original
declaration order — so the body's indices always line up with what's
actually built.
- `HigherOrderFunctionExpr::evaluate` forwards `lambda.used_params()` to
`LambdaArgument::new`

## Are these changes tested?

yes, added two new tests one for the unused-parameter case and
nested-lambda for the shadowing case.

## Are there any user-facing changes?

The only public api change is on `LambdaArgument::new ` which now
requires a new argument: `used_params: &HashSet<String>`, however
LambdaArgument::new is very unlikely to be called outside datafusion,
see
[this](apache#22853 (comment))
comment

(cherry picked from commit 4e6acfe)

* Adjust to API change
LiaCastaneda added a commit that referenced this pull request Aug 14, 2026
…Partition child (#169)

* Revert "fix(lambda): only push referenced params into the merged batch (apache#24162) (#166)"

This reverts commit 79de5e9.

* fix: keep a CoalescePartitionsExec required by a SinglePartition child (apache#23948)

- None filed; happy to open one if preferred.

A valid query can be planned into a physical plan that `SanityCheckPlan`
then rejects:

```
SanityCheckPlan
caused by
Error during planning: Plan: ["HashJoinExec: mode=CollectLeft, join_type=Left, on=[(id@0, id@0)], projection=[id@0]",
  "  DataSourceExec: file_groups={4 groups: [...]}, projection=[id], file_type=parquet",
  "  RepartitionExec: partitioning=RoundRobinBatch(8), input_partitions=1",
  "    CoalescePartitionsExec",
  "      ProjectionExec: expr=[first_value(t.id) ORDER BY [...]@1 as id]",
  "        AggregateExec: mode=FinalPartitioned, gby=[id@0 as id], aggr=[first_value(t.id) ORDER BY [...]]",
  "          RepartitionExec: partitioning=Hash([id@0], 8), input_partitions=4",
  "            AggregateExec: mode=Partial, gby=[id@1 as id], aggr=[first_value(t.id) ORDER BY [...]]",
  "              DataSourceExec: file_groups={4 groups: [...]}, projection=[ts, id], file_type=parquet"]
does not satisfy distribution requirements: SinglePartition. Child-0 output partitioning: UnknownPartitioning(4)
```

The `HashJoinExec` is in `CollectLeft` mode, which requires
`Distribution::SinglePartition` on its build (left) child, but child 0
is a bare 4-partition `DataSourceExec` with no `CoalescePartitionsExec`
above it.

Self-contained reproducer with `datafusion-cli` (the four `COPY`
statements are what make the scan multi-partition):

```sql
set datafusion.execution.target_partitions = 8;
set datafusion.optimizer.repartition_file_scans = false;

create table src (id int, ts int) as values (1, 10), (2, 20), (3, 30);

copy (select * from src) to 'data/0.parquet' stored as parquet;
copy (select * from src) to 'data/1.parquet' stored as parquet;
copy (select * from src) to 'data/2.parquet' stored as parquet;
copy (select * from src) to 'data/3.parquet' stored as parquet;

create external table t stored as parquet location 'data/';

select a.id
from t a
left join (select distinct on (id) id, ts from t order by id, ts) f on a.id = f.id
order by a.id;
```

Setting `datafusion.optimizer.repartition_sorts = false` makes it plan
fine, which points at the sort-parallelization phase.

`EnsureRequirements` does insert the coalesce for the `SinglePartition`
requirement (`enforce_distribution.rs`, `Distribution::SinglePartition
=> add_merge_on_top(...)`). Its own phase 3a (`parallelize_sorts`) then
takes it back out: `remove_bottleneck_in_subplan` removes a
`CoalescePartitionsExec` found at `children[0]` positionally, without
consulting the parent's distribution requirement for that child.

That parent is reached because `update_coalesce_ctx_children` marks a
node as connected when *any* child qualifies. It correctly excludes a
`SinglePartition`-requiring child from *setting* the flag, but the
join's other child (`UnspecifiedDistribution`, connected to a coalesce
below) sets it, so the traversal descends into the join and rewrites
child 0 anyway. Nothing re-enforces distribution afterwards, so
`SanityCheckPlan` is the first thing to notice. Note the surviving
`CoalescePartitionsExec` on the probe side in the plan above: it is what
propagated the flag, and it is untouched because the `if` returns
without recursing into child 1.

The sibling helper on the phase 2b path already does consult the
requirement (`update_child_to_remove_unnecessary_sort` /
`remove_corresponding_sort_from_sub_plan` re-add a merge using the
per-child `child_distribution(child_idx)`); only this path is missing
it.

The same failure shows up with a build child that is already
hash-partitioned on the join key (`Child-0 output partitioning:
Hash([k@0], 8)`), which is what a `JoinSelection` input swap leaves
behind — a `CollectLeft` join reported as `join_type=Right` with an
embedded projection.

`remove_bottleneck_in_subplan` now checks the parent's per-child
distribution requirement before removing a coalesce, both for
`children[0]` and when recursing into the other children.

The node `parallelize_sorts` is itself rewriting (the root of the call)
is exempt, since the caller drops that node and rebuilds the sort
cascade around the result — that is the rule's intended transformation,
and gating it too would disable sort parallelization below a global
sort. This is threaded through as an `is_root` flag on a private `_impl`
function; the public entry point keeps its signature.

Yes, at two levels:

- An end-to-end sqllogictest in
`datafusion/sqllogictest/test_files/joins.slt` reproducing it from SQL
(the reproducer above, with the data written by `COPY` inside the test).
On `main` it fails with exactly the distribution error above.
- Two tests in
`datafusion/core/tests/physical_optimizer/ensure_requirements.rs`
covering both shapes of the build child (`UnknownPartitioning(n)` and
`Hash([k], n)`), running the full `EnsureRequirements` rule and then
`SanityCheckPlan` via the existing `optimize_and_sanity_check` helper,
plus the idempotency check.

`cargo test -p datafusion-physical-optimizer`, `cargo test -p datafusion
--test core_integration -- physical_optimizer` (530 tests) and the full
`sqllogictest` suite (498 files) pass.

No API changes. Plans that were previously rejected by `SanityCheckPlan`
now plan and execute; a coalesce that is genuinely required is retained
where it was previously (incorrectly) removed.

---------

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>

* fix(lambda): only push referenced params into the merged batch (apache#24162) (#166)

* fix(lambda): only push referenced params into the merged batch (apache#24162)

## Which issue does this PR close?

basically this PR apache#22853 + a
few more tests

## Rationale for this change

The current lambdas in DF only take a single parameter `(v -> ...)`, so
nobody had noticed that `LambdaExpr` mishandles lambdas with more than
one parameter. The bug surfaced while working on `transform_values`
(apache#22689), which needs `(k, v) -> expr ` two parameters, one of which is
very often unused (e.g. `(k, v) -> v * 2`, k never referenced).

The bug is that when a higher order function with more than 1 param
evaluates a lambda, it fills each parameter into a slot based on its
declared position — for example for `(k, v) -> v` `k` always goes into
slot 0, `v` always into slot 1. `LambdaExpr` separately scans the body
and renumbers whatever it finds referenced into a dense `0..n` range, to
avoid carrying around columns nothing uses (like `v` in this case). That
renumbering is fine for outer captures, but applying it to the lambda's
own parameters is wrong, because it changes where the body looks for a
value without changing where the evaluator put it.

### Example:

in `(k, v) -> v` `v` is declared second (slot 1), but since it's the
only parameter the body references, the renumbering logic reassigns it
to slot 0. The evaluator, unaware of this, writes `k`'s values into slot
0 and `v`'s into slot 1. So the body ends up reading slot 0 expecting
`v` — and gets `k` instead. So the results end up being incorrect.

## What changes are included in this PR?

- `LambdaExpr` now computes `used_params`: which is the subset of its
own declared parameters that are actually referenced in the body.
- `LambdaArgument::new` takes `used_params` and only pushes the
referenced parameters in the body into the merged batch, in original
declaration order — so the body's indices always line up with what's
actually built.
- `HigherOrderFunctionExpr::evaluate` forwards `lambda.used_params()` to
`LambdaArgument::new`

## Are these changes tested?

yes, added two new tests one for the unused-parameter case and
nested-lambda for the shadowing case.

## Are there any user-facing changes?

The only public api change is on `LambdaArgument::new ` which now
requires a new argument: `used_params: &HashSet<String>`, however
LambdaArgument::new is very unlikely to be called outside datafusion,
see
[this](apache#22853 (comment))
comment

(cherry picked from commit 4e6acfe)

* Adjust to API change

---------

Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com>
Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
Co-authored-by: Lía Adriana <lia.castaneda@datadoghq.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants