Skip to content

Cherry-pick apache/datafusion#25185 - #182

Closed
LiaCastaneda wants to merge 81 commits into
DataDog:branch-55from
LiaCastaneda:lia.castaneda/cherry-pick/apache-pr-25185-20260916
Closed

LiaCastaneda wants to merge 81 commits into
DataDog:branch-55from
LiaCastaneda:lia.castaneda/cherry-pick/apache-pr-25185-20260916

Conversation

@LiaCastaneda

Copy link
Copy Markdown

cherry-picks apache#25185

Also fixes a cache-invalidation gap in take_n exposed by conflict resolution with DataDog's pre-existing partial port of this optimization (#171): cached_values was not cleared after take_n remaps inner slots, so a stale cache could be reused against remapped data. Added self.cached_values = None; before the rebuild, matching the doc comment's stated invariant. All existing dictionary group-values tests (13/13) plus the broader aggregates suite (155/155) pass.

mbutrovich and others added 30 commits May 20, 2026 16:37
…20337)" (apache#22437) (apache#22445)

- Backports apache#22437 from @alamb
to the branch-54 line

This PR cherry-picks the revert of `ExecutionPlan::apply_expressions()`
(apache#20337) onto `branch-54` so that DataFusion 54.0 does not ship the new
public API.
…apache#22443)

## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

Backport apache#22404 

- Closes #.

## Rationale for this change

<!--
Why are you proposing this change? If this is already explained clearly
in the issue then this section is not needed.
Explaining clearly why changes are proposed helps reviewers understand
your changes and offer better suggestions for fixes.
-->

## What changes are included in this PR?

<!--
There is no need to duplicate the description in the issue here but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->

## Are these changes tested?

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->

## Are there any user-facing changes?

<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.
-->

<!--
If there are any breaking changes to public APIs, please add the `api
change` label.
-->
…ter stats-based file reorder (apache#22501)

## Which issue does this PR close?

Cherry-pick of apache#22493 onto `branch-54`.

## Rationale for this change

`branch-54` includes apache#21956 (`feat: globally reorder files and row
groups by statistics for TopK queries`), which introduced a regression:
for plain-column, multi-file scans where the on-disk file order does not
match the declared sort order, `SortExec` was no longer eliminated even
when stats-based reorder produced non-overlapping file groups whose
declared ordering re-validated.

apache#22493 restores the pre-apache#21956 sort-elimination behaviour by
re-validating `output_ordering` after `rebuild_with_source` reorders
files, and (per @adriangb's correctness follow-up) restoring the
original hint-free `file_source` on the Inexact→Exact upgrade so
leftover `reverse_row_groups` / `sort_order_for_reorder` hints don't
mis-order row groups within a single file once the `SortExec` safety net
is gone.

## What changes are included in this PR?

Straight cherry-pick of merge commit `94c58d086`. Includes:

- `FileScanConfig::try_pushdown_sort` Inexact arm: re-validate, upgrade
to Exact (with file_source restore), guard with NULL safety +
early-return
- `rebuild_with_source`: `match (all_non_overlapping, is_exact)`
decision table for keep_ordering
- SLT updates restoring `SortExec` elimination expectations + Tests
5b/5c/8b for the NULL-safety and same-min row-group edge cases

## Are these changes tested?

Cherry-picked cleanly (auto-merge in `sort_pushdown.rs`).
`cargo build -p datafusion-datasource` — passes.
`cargo test -p datafusion-sqllogictest --test sqllogictests --
sort_pushdown` — passes.

## Are there any user-facing changes?

Same as apache#22493: plain-column wrong-order-files cases regain SortExec
elimination when files happen to be non-overlapping by statistics. No
new API.

Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com>
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
… container types (apache#21934) (apache#22446)

- Backports apache#21934 from
@bert-beyondloops to the branch-54 line

This PR cherry-picks the fix for `ScalarValue::compact` to compact view
buffers for all container types onto `branch-54`.

Co-authored-by: Bert Vermeiren <103956021+bert-beyondloops@users.noreply.github.com>
Co-authored-by: Bert Vermeiren <bert.vermeiren@datadobi.com>
Co-authored-by: Dmitrii Blaginin <dmitrii@blaginin.me>
## Which issue does this PR close?

- Refs apache#22557.
- Companion PR to apache#22559, targeting `branch-54`.

## Rationale for this change

DataFusion 54 changed `ExecutionPlan` downcasting to use the `Any`
supertrait directly. That removes `ExecutionPlan::as_any`, which had
also served as a customization point for wrapper nodes: wrappers could
identify as themselves internally while exposing the wrapped plan type
to normal downcast-based inspection.

This PR adds an explicit `ExecutionPlan::downcast_delegate()` hook for
wrapper nodes that want their public `ExecutionPlan` downcast identity
to be delegated to another plan.

The proposed behavior intentionally preserves the old `as_any` override
semantics: when a node opts into downcast delegation, intermediate
delegating wrappers are invisible to `dyn ExecutionPlan::is::<T>()` and
`downcast_ref::<T>()`.

## What changes are included in this PR?

- Adds `ExecutionPlan::downcast_delegate()` with a default
implementation returning `None`.
- Updates `dyn ExecutionPlan::is::<T>()` and `downcast_ref::<T>()` to
delegate to `downcast_delegate()` when present, otherwise use the
current concrete plan type.
- Documents that `downcast_delegate()` is only for type introspection
and is independent from `children()` / plan traversal.
- Adds tests for direct and nested downcast-delegating wrappers,
including that intermediate delegating wrappers remain invisible to
normal downcast-based inspection.

## Are these changes tested?

Yes.

- `cargo test -p datafusion-physical-plan execution_plan_downcast`
- `cargo test -p datafusion-physical-plan --lib`
- `cargo fmt --all -- --check`
- `git diff --check`

## Are there any user-facing changes?

Yes. This adds a new public `ExecutionPlan` trait method with a default
implementation, and it changes `ExecutionPlan` downcast helpers to honor
wrappers that explicitly opt into delegating public downcast identity.
) (apache#22634)

- Part of apache#21080

This PR:
- Backports apache#22571 from
@kumarUjjawal to the `branch-54` line

Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com>
…erUDF struct (apache#22593) (apache#22635)

- Part of apache#21080

This PR:
- Backports apache#22593 from
@LiaCastaneda to the `branch-54` line

## Note on conflict resolution

The cherry-pick had conflicts in two test areas:
- `datafusion/functions-nested/src/array_any_match.rs` — only the
test-module `use` line conflicted; `branch-54`'s test does not reference
the renamed symbols, so the existing import was kept (the production
`HigherOrderUDF` -> `HigherOrderUDFImpl` rename applied cleanly).
- `datafusion/substrait/tests/cases/roundtrip_logical_plan.rs` — the
upstream PR modified the
`roundtrip_array_transform_higher_order_function` test and
`ArrayTransform` helper, but that test was added to `main` after
`branch-54` was cut and does not exist on `branch-54`. It is out of
scope for this refactor backport, so the `branch-54` file was left
unchanged.

Co-authored-by: Lía Adriana <lia.castaneda@datadoghq.com>
…ryToJoin` (#… (apache#22693)

This is PR packports apache#22316
from @neilconway to branch-54

Co-authored-by: Neil Conway <neil.conway@gmail.com>
…pache#22530) (apache#22690)

## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- Closes #.

## Rationale for this change

This PR backports apache#22530


## What changes are included in this PR?

<!--
There is no need to duplicate the description in the issue here but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->

## Are these changes tested?

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->

## Are there any user-facing changes?

<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.
-->

<!--
If there are any breaking changes to public APIs, please add the `api
change` label.
-->
…Pushdown (apache#22525) (apache#22631)

- Part of apache#21080

This PR:
- Backports apache#22525 from
@kumarUjjawal to the `branch-54` line

Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com>
…flag (backport apache#22632) (apache#22648)

## Which issue does this PR close?

- Backport of apache#22632 to `branch-54`.

## Rationale for this change

Content-defined chunking (CDC) write options were added in apache#21110 and
are slated for the 54.0.0 release. This backports the refactor in apache#22632
so the config/proto surface ships in its final form, before the release
goes out.

The CDC options previously worked as `use_content_defined_chunking:
Option<CdcOptions>` with a `ConfigField` impl that accepted a bare
`use_content_defined_chunking = true|false` and otherwise enabled CDC
implicitly when any sub-field was set. This has a few problems:

- **Naming diverges from parquet-rs.** `WriterProperties` exposes
`content_defined_chunking()` /
`set_content_defined_chunking(Option<CdcOptions>)` with no `use_`
prefix.
- **Implicit / order-dependent on the SQL side.** Format options in
`COPY ... OPTIONS` / `CREATE EXTERNAL TABLE ... OPTIONS` are applied
from a `HashMap` (non-deterministic order). With the old bare-boolean
form, mixing `... = false` with a sub-field could resolve to enabled or
disabled depending on iteration order.
- **Extra machinery.** Supporting the bare boolean required hand-written
`ConfigField` impls and a `#[expect(clippy::should_implement_trait)]`
workaround, plus a zero-sentinel fallback in the proto mapping.

Since CDC is unreleased, the config/proto surface can still be changed
freely.

## What changes are included in this PR?

- Rename the `ParquetOptions` field `use_content_defined_chunking` ->
`content_defined_chunking` (matches parquet-rs).
- Make `CdcOptions` a plain `config_namespace!` with an explicit
`enabled: bool` field alongside the chunking parameters; the field is a
bare `CdcOptions` (no longer `Option<CdcOptions>`). CDC is on iff
`content_defined_chunking.enabled` is true. Setting a parameter no
longer implicitly enables CDC, and the result is independent of key
order.
- Add `CdcOptions::enabled()` / `CdcOptions::disabled()` shorthand
constructors.
- Drop the `ConfigField` impls and the `should_implement_trait`
workaround — all generated by the macro now.
- Add an `enabled` field to the proto `CdcOptions` message so the proto
<-> config mapping is a plain field copy in both directions.
- Update unit tests, regenerate config docs + the `information_schema`
snapshot, and add `parquet_cdc_config.slt` documenting the resolution
behavior.

## Are these changes tested?

Yes — `datafusion-common` config + writer unit tests,
`datafusion-proto-common` proto round-trip tests, `datafusion/core`
parquet integration tests, and sqllogictest (`parquet_cdc.slt` + new
`parquet_cdc_config.slt`). Cherry-pick applied cleanly onto `branch-54`;
affected crates build and the CDC unit tests pass.

## Are there any user-facing changes?

Yes, but only to the unreleased CDC options:
- Config key `datafusion.execution.parquet.use_content_defined_chunking`
-> `datafusion.execution.parquet.content_defined_chunking.enabled` (plus
`.min_chunk_size` / `.max_chunk_size` / `.norm_level`).
- The bare-boolean form is removed; enable/disable via
`content_defined_chunking.enabled = true|false`.

No released API is affected.

🤖 Generated with [Claude Code](https://claude.com/claude-code)
## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- Closes #.

## Rationale for this change

<!--
Why are you proposing this change? If this is already explained clearly
in the issue then this section is not needed.
Explaining clearly why changes are proposed helps reviewers understand
your changes and offer better suggestions for fixes.
-->

## What changes are included in this PR?

<!--
There is no need to duplicate the description in the issue here but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->

## Are these changes tested?

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->

## Are there any user-facing changes?

<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.
-->

<!--
If there are any breaking changes to public APIs, please add the `api
change` label.
-->

---------

Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
…trip (backport apache#22104) (apache#22785)

## Which issue does this PR close?

- Backport of apache#22104 to `branch-54` (for 54.1.0, tracked in apache#22547).

This PR:
- Backports apache#22104 to the `branch-54` line so the `null_aware` proto
  round-trip fix ships in 54.1.0, as requested in

apache#22065 (comment)

Clean cherry-pick; `datafusion-proto` builds and both round-trip
regression tests pass on `branch-54`.
apache#22453) (apache#126)

## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- Closes #.

## Rationale for this change

When the substrait consumer hits an `Aggregate` with two identical
measures (e.g. `sum(a)` present twice), planning fails with `Schema
contains duplicate unqualified field name`. Substrait carries column
names at the plan root rather than on the measures themselves, so the
measures arrive at `Aggregate` schema construction without aliases --
and two identical exprs produce two identical field names. PR apache#20539
fixed the `NameTracker` to dedupe duplicate names in the consumer, but
it was only applied to grouping expressions, not to the measures.

The planner sees:

```
field 1: (qualifier: None, name: "sum(data.a)")
field 2: (qualifier: None, name: "sum(data.a)")
```

which is rejected when constructing the Aggregate's output schema.

## What changes are included in this PR?

Run aggregate measures through the same `NameTracker` like the grouping
expressions in `from_aggregate_rel`

## Are these changes tested?

Yes -- added a roundtrip test `aggregate_identical_measures`. Without
the fix it produces `Error: SchemaError(DuplicateUnqualifiedField {
name: "sum(data.a)" }, Some(""))`

## Are there any user-facing changes?

No.

(cherry picked from commit 097efae)
Part of apache#21172

Substrait support wasn't implemented in the core lambda support to
reduce PR size

Substrait consuming and producing of higher-order functions, lambdas and
lambda variables

Unit tests added to
`datafusion/substrait/tests/cases/roundtrip_logical_plan.rs`

None

---------

(cherry picked from commit 9a6f67e)
(cherry picked from commit 1ac2df1)

Co-authored-by: gstvg <28798827+gstvg@users.noreply.github.com>
Co-authored-by: Raz Luvaton <16746759+rluvaton@users.noreply.github.com>
Co-authored-by: Ben Bellick <36523439+benbellick@users.noreply.github.com>
…atch (apache#22852)

## Which issue does this PR close?

- Closes apache#22849
- A related cross-partition starvation case is tracked separately in
apache#22874 and addressed by an upcoming follow-up PR — see
[discussion](apache#22852 (comment))
for details

## Rationale for this change

`TopK::insert_batch` short-circuits when the heap's dynamic filter
rejects every row in a batch:

```rust
if !filter.has_true() {
    // nothing to filter, so no need to update
    return Ok(());
}
```

The early-exit check `attempt_early_completion(&batch)` lives later in
the same function, gated on `replacements > 0`. So a batch that the
filter rejects entirely bypasses the check.

The heap's dynamic filter is derived from the heap's worst row (via
`update_filter`). A batch whose rows all come from a strictly worse sort
prefix is exactly the batch the filter rejects entirely — i.e. the very
signal `attempt_early_completion` is designed to detect ("the next batch
is past the heap's boundary, we can stop") is what causes the function
to short-circuit *before* the check runs.

This is a feature-interaction regression between two PRs that were both
correct in isolation. The `attempt_early_completion` mechanism was added
by apache#15563 (closing apache#15529). At the time, there was no heap-derived
dynamic filter on TopK, so the only sensible call site was right after a
successful heap insertion. Two months later, apache#15770 added the
dynamic-filter pushdown for TopK sorts, introducing the
`!filter.has_true()` short-circuit. The two features address different
problems and the new short-circuit didn't connect to the existing
prefix-completion check — which is how this gap opened up.

**Consequence**: on a TopK over an input ordered on the sort prefix,
`finished = true` is never set once the heap stabilizes. Since
`finished` is the signal `SortExec` uses to stop pulling from its input
(via `Poll::Ready(None)` from the TopK stream, which cascades into
dropping the source stream), the source keeps being polled long past the
point where no further row can improve the heap. The LIMIT optimization
effectively degrades to "heap saves memory but reads everything";
sources with cancellable streams (e.g. networked sources) never receive
the cancellation signal.

## What changes are included in this PR?

Single behavioral change in `datafusion/physical-plan/src/topk/mod.rs`:
call `attempt_early_completion(&batch)` immediately before the `return
Ok(())` in the `!filter.has_true()` branch.

Why this scope, not a broader restructuring:

- The existing `attempt_early_completion` call inside `if replacements >
0` is load-bearing for a related case: a batch containing a mix of
"still valuable" rows and "past the boundary" rows. The existing
`test_try_finish_marks_finished_with_prefix` test covers this case —
Batch 2 with `a=[2,3], b=[10,20]` against a heap where `heap.max.a = 2`;
the `(2, 10)` row must be inserted before the check on the `(3, 20)`
last row triggers. Moving the call earlier would skip the insertion of
valuable rows and break that test.
- The bug is specifically that the *short-circuit* path doesn't call the
check. The fix targets exactly that path.
- A related but separate gap is not addressed here: when
`filter.has_true() == true` but `replacements == 0` (the filter accepts
some rows but `find_new_topk_items` ends up inserting none of them), the
existing call inside `if replacements > 0` is also skipped. This
requires a divergence between the heap's filter predicate and the
row-byte comparison used inside `find_new_topk_items`, which shouldn't
normally happen (the filter is derived from the heap's worst row using
the same comparator). A deterministic synthetic repro would likely
require concurrent heap updates from sibling partitions or
boundary-value edge cases (NaN/NULL semantics, type coercion). Happy to
send a follow-up if reviewers want it covered; the workload that
motivated this fix was the filter-rejection case empirically.

## Are these changes tested?

Yes. Added a regression test
`test_try_finish_fires_when_filter_rejects_entire_batch`. The assertion
target is `topk.finished` — the flag that signals "stop pulling from the
source" to upstream consumers (read by `TopKExec::poll_next` to emit
`Poll::Ready(None)`). Asserting that the flag transitions on the
fully-filter-rejected batch is equivalent to asserting that the
source-stopping mechanism activates.

- Builds a TopK over a `(a, b)` sort with prefix `a`, k=3.
- Inserts a batch that fills the heap with rows from `a ∈ {1, 2}`;
`update_filter` tightens the filter to `a < 2 OR (a = 2 AND b < 30)`.
- Inserts a second batch with all rows at `a = 3` — filter rejects every
row.
- Without the fix: `insert_batch` short-circuits, `topk.finished` stays
`false`. Test fails.
- With the fix: `attempt_early_completion` fires (last-row prefix `a =
3` > heap.max prefix `a = 2`), `topk.finished` becomes `true`. Test
passes.

The test also asserts the emitted top-K is unchanged from after batch 1,
confirming no candidate row was incorrectly excluded by the early bail.

All 28 existing `topk::` tests continue to pass (including
`test_try_finish_marks_finished_with_prefix`, which exercises the
mixed-prefix case).

## Are there any user-facing changes?

No public API or output changes. The fix only changes when TopK marks
itself `finished = true` — specifically, it now fires
`attempt_early_completion` for batches that are entirely rejected by the
heap's dynamic filter, where previously it would silently skip the
check. Output of TopK is unchanged; only the early-exit behavior
improves.

---------

Co-authored-by: Gabriel <45515538+gabotechs@users.noreply.github.com>
(cherry picked from commit 6520315)
…ache-pr-22852-branch54-20260617

[branch-54] Cherry-pick apache#22852

Co-authored-by: ajegou <arnaud.jegou@gmail.com>
Co-authored-by: arnaud.jegou <arnaud.jegou@datadoghq.com>
Physical plan proto serialization was still binding the plan as dyn Any after the as_any removal. That bypassed ExecutionPlan::downcast_ref, so wrapper plans that delegate their public downcast identity fell through to the extension codec instead of serializing as the wrapped built-in plan.

Bind the serializer view as dyn ExecutionPlan so the existing downcast chain uses the delegating helper, and add a regression test with a wrapper around EmptyExec.

(cherry picked from commit 401d8fc)
…cast-delegate

fix(proto): honor ExecutionPlan downcast_delegate during serialization
Cherry pick a0763db - Fix ArrayCompact incompatibility
Cherry pick 6692f6f - fix(substrait): dedupe names
mattp5657 and others added 26 commits August 7, 2026 14:40
…apache#23583)

## Which issue does this PR close?

- Closes apache#23454.

## Rationale for this change

apache#23184 let compatible range-partitioned inputs satisfy inner partitioned
hash joins without repartitioning. Full partitioned equi joins still
always went through the conservative hash-repartition path, even when
both inputs were already co-partitioned by range on the join key(s). The
per-partition unmatched-row tracking in `HashJoinExec` is already
partition-local under `PartitionMode::Partitioned` (not shared globally
like in `CollectLeft`), so Full-join semantics generalize cleanly to
range co-partitioning with no additional bookkeeping required.


## What changes are included in this PR?

- Extend `HashJoinExec::input_distribution_requirements()` to opt
`JoinType::Full` in to
`allow_range_satisfaction_for_key_partitioning()`, alongside the
existing `JoinType::Inner` case. The underlying `co_partitioned` /
`compatible_co_partitioning_layout` / `co_partitioning_satisfied` logic
in `distribution_requirements.rs` was already join-type-agnostic, so no
changes were needed there.
- Add planner unit tests in `enforce_distribution.rs` covering both the
compatible-layout case (no repartition inserted) and the
incompatible-split-points case (repartition still inserted) for
`JoinType::Full`.
- Add a `range_partitioned_sparse` sqllogictest fixture table with the
same partition layout as `range_partitioned` but only partially
overlapping keys, and add sqllogictest coverage in
`range_partitioning.slt` for: a compatible Full join avoiding
repartition, an incompatible Full join still repartitioning, and
matched/left-only/right-only unmatched rows produced correctly by a
co-partitioned Full join.

## Are these changes tested?

Yes:
- Two new Rust unit tests in
`datafusion/core/tests/physical_optimizer/enforce_distribution.rs`
assert on the physical plan shape (repartition inserted or not) for
compatible and incompatible range layouts.
- Three new sqllogictest cases in
`datafusion/sqllogictest/test_files/range_partitioning.slt` exercise the
feature end-to-end against real data, including matched rows, left-only
unmatched rows, and right-only unmatched rows for a Full outer join.
- Existing `enforce_distribution` tests, `range_partitioning.slt`, and
proto roundtrip tests all continue to pass.


## Are there any user-facing changes?

Yes: `EXPLAIN` output for `Full` joins over compatible range-partitioned
inputs will no longer show a `RepartitionExec`, and such queries will
avoid the associated hash-shuffle cost at execution time. No public API
changes.

---------

Co-authored-by: Matthew Patton <matthewpatton@macbookpro.mynetworksettings.com>
…#23484)

## Which issue does this PR close?

- Closes apache#23453.
- Part of apache#22395.

## Rationale for this change

apache#23184 let compatible range-partitioned inputs satisfy **inner**
partitioned hash joins without repartitioning. Right-side partitioned
equi joins have the same locality guarantee: `Right`, `RightSemi`,
`RightAnti` and `RightMark` all anchor every output row on the probe
(right) partition (`on_lr_is_preserved` is probe-side for all four), so
when both inputs are co-partitioned by the join keys, execution stays
partition-local and the hash repartition is unnecessary.

This removes unnecessary `RepartitionExec`s for already-co-located
inputs, extending the inner-join behavior from apache#23184 to the right-side
variants. No micro-benchmark included, consistent with apache#23184;
correctness is demonstrated by matched/unmatched execution results
below.

## What changes are included in this PR?

The only production change is widening the existing inner-only gate in
`HashJoinExec::input_distribution_requirements` from `join_type ==
JoinType::Inner` to `matches!(join_type, Inner | Right | RightSemi |
RightAnti | RightMark)` (under `PartitionMode::Partitioned`). The
`co_partitioned` / range-satisfaction machinery from apache#23184 is unchanged
— the sanity checker and enforce_distribution consume range satisfaction
join-type-agnostically.

One behavior note for apache#23376: range co-partitioned right joins now stay
`Range`/`Range`, so partitioned dynamic filters are disabled for them
(`has_partitioned_dynamic_filter_routing` returns false), the same safe
delta apache#23184 introduced for inner joins. These variants are
probe-not-preserved for pruning, so dynamic filters were not eligible to
prune their probe rows regardless.

## Are these changes tested?

Yes.

- Optimizer (`enforce_distribution.rs`): 4 reuse tests (one per join
type) proving compatible range/range inputs keep `Range` partitioning
with no `RepartitionExec`; 4 incompatibility tests proving that
mismatched split points, sort options, partition counts, or join-key
expressions still insert a hash repartition; sanity-check pairs
mirroring apache#23184.
- Execution (`range_partitioning.slt`): `EXPLAIN` plan pins plus
matched/unmatched result checks for `Right` (incl. `NULL` left values),
`RightSemi`, `RightAnti`, and an incompatible-layout `Right` join
(repartitions, correct results).
- The new reuse tests fail on `main` without the gate change (verified
by reverting the production diff: exactly the 4 reuse tests fail,
everything else passes).
- `RightMark` testing note: `RightMark` is not reachable from SQL in
sqllogictest (`IN`-subquery decorrelation emits `LeftMark`; physical
`RightMark` only appears via a statistics-based swap), so its
co-partitioning behavior is pinned at the optimizer/plan level, and mark
null-marker semantics (matched/unmatched/`NULL` build keys) are pinned
via the `LeftMark` path in the slt.

## Are there any user-facing changes?

No API changes. Plans over compatible range-partitioned inputs avoid a
hash repartition for right-side equi hash joins.
## Which issue does this PR close?

- Closes apache#23230 

## Rationale for this change

After apache#23231 was merged in for supporting physical execution of the
range repartitioning scheme, we still had a few methods on physical
planning unimplemented, specifically `try_swapping_with_projection`,
`try_pushdown_sort`, `repartitioned` - this PR finishes the
implementation of those methods

## What changes are included in this PR?

- `try_swapping_with_projection`: similar to the `Hash` scheme, for
`Range` we call `update_expr` for each of the range key expressions to
attempt rewriting based on the projection expressions
- `try_pushdown_sort`: same as other variants, we delegate to the child
and wrap with a new `RepartitionExec`
- `repartitioned`: unable to support for Range, left comment in codebase
with explanation

## Are these changes tested?

Yes

## Are there any user-facing changes?

No
…for left-side hash joins (apache#23487)

## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- Closes apache#23452 

## Rationale for this change

Allows compatible `Partitioning::Range` inputs to satisfy partitioned
hash join distribution requirements for left-side joins, avoiding
unnecessary hash repartitioning.

<!--
Why are you proposing this change? If this is already explained clearly
in the issue then this section is not needed.
Explaining clearly why changes are proposed helps reviewers understand
your changes and offer better suggestions for fixes.
-->

## What changes are included in this PR?

- Enables range co-partitioning satisfaction for `Left`, `LeftSemi`,
`LeftAnti`, and `LeftMark` hash joins.
- Adds optimizer and sqllogictest coverage for compatible and
incompatible range layouts.
- Covers matched/unmatched rows and LeftMark null-related marker
behavior.

<!--
There is no need to duplicate the description in the issue here but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->

## Are these changes tested?

Yes. Added/updated physical optimizer tests and
`range_partitioning.slt`.

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->

## Are there any user-facing changes?

<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.
-->

<!--
If there are any breaking changes to public APIs, please add the `api
change` label.
-->
## Which issue does this PR close?

- Closes apache#22394 

## Rationale for this change

Exposing range partition metadata via the FFI for external consumers.

## What changes are included in this PR?

- Added FFI mirror struct for `RangePartitioning` and added new enum
variant for range in `FFI_Partitioning`
- For native -> FFI, added match arm for the new variant, same with FFI
-> native but changed the approach of `From` -> `TryFrom` to utilize the
validation for `RangePartitioning` and modified `plan_properties` to
match
- Added tests


## Are these changes tested?

Yes

## Are there any user-facing changes?
Yes, exposing Range partitioning over FFI. This exposes a new `Range`
variant in the `FFI_Partitioning` enum, which may cause consumers of
this enum to add another arm to match statements to handle the new enum.
New `FFI_RangePartitioning` struct for the `Range` variant.

---------

Co-authored-by: Tim Saucer <timsaucer@gmail.com>
<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- Closes apache#23266.
- Part of apache#22395.

<!--
Why are you proposing this change? If this is already explained clearly
in the issue then this section is not needed.
Explaining clearly why changes are proposed helps reviewers understand
your changes and offer better suggestions for fixes.
-->

`Partitioning::Range` now satisfies `Distribution::KeyPartitioned`
privately across all operators that require it: aggregates, windows,
TopK, and co-partitioned joins. So now the temporary operator opt-ins
can be consolidated.

<!--
There is no need to duplicate the description in the issue here but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->

- Allow compatible `Partitioning::Range` to satisfy
`Distribution::KeyPartitioned` via `Partitioning::satisfaction`
- Remove the temporary range-satisfaction helpers and operator-specific
opt-ins

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->

Yes

<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.
-->

No, this should not change any exisitng behavior just consolidation

<!--
If there are any breaking changes to public APIs, please add the `api
change` label.
-->
…ng-cherry-pick

Branch 54 range partitioning cherry pick

Co-authored-by: gene-bordegaray <gene.bordegaray@datadoghq.com>
Co-authored-by: saadtajwar <59696464+saadtajwar@users.noreply.github.com>
Co-authored-by: JSOD11 <justin.odwyer@datadoghq.com>
Co-authored-by: gmhelmold <gustavomalleths@gmail.com>
Co-authored-by: mattp5657 <matthewpatton641@gmail.com>
Co-authored-by: Rich-T-kid <137434454+Rich-T-kid@users.noreply.github.com>
Co-authored-by: EdsonPetry <124717297+EdsonPetry@users.noreply.github.com>
Co-authored-by: stuhood <stuhood@gmail.com>
Co-authored-by: mithuncy <mithun.cy@gmail.com>
…e#24162) (apache#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
…Partition child (apache#169)

* Revert "fix(lambda): only push referenced params into the merged batch (apache#24162) (apache#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) (apache#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>
…mas on the column-wise path (apache#23523)

- Part of apache#22715 (nested type coverage in `GroupValuesColumn` EPIC)
- Alternative to apache#23128 (per-type approach) — implements the direction
@alamb proposed in
<apache#23128 (comment)>
- Step toward the terminal goal of retiring `GroupValuesRows` entirely
(apache#23404)

Today `GroupValuesColumn` is **all-or-nothing**: a single nested column
in the GROUP BY key (\`Struct\`, \`List\`, \`FixedSizeList\`, …) makes
\`supported_schema\` return \`false\` and drops the *entire* aggregation
onto the row-wise \`GroupValuesRows\` fallback — even when every other
column would have qualified for the column-wise fast path. For a \`GROUP
BY int_col, struct_col\` shape, the \`int_col\` pays the row-encoded
storage cost for no reason.

Add \`RowsGroupColumn\`: a generic \`GroupColumn\` backed by a
single-field \`RowConverter\`, wired in as the nested-type dispatch arm
of \`group_column_supported_type\` / \`make_group_column\`. Native
columns keep their type-specialized builders; the nested column pays
row-encoding only for its one column.

Gated to \`data_type.is_nested()\` so intentionally excluded scalar
types (Float16, Decimal256) stay on \`GroupValuesRows\` and the
\`group_column_supported_type\` ⇔ \`make_group_column\` invariant holds.

Memory, measured with 4000 groups of \`8 × Int64 + 1 ×
FixedSizeList<Int64, 4>\` in
\`mixed_schema_column_path_uses_less_memory_than_rows_fallback\`:

| | Bytes | vs baseline |

|-------------------------------------------------------|----------|-------------|
| \`GroupValuesRows\` (today's fallback) | 1096 KB | 100% |
| \`GroupValuesColumn\` + \`RowsGroupColumn\` fallback | 594 KB |
**54.2%** |

Speed: not benchmarked as a headline result — the wins come from native
columns keeping their type-specialized \`equal_to\`/\`append_val\` fast
paths instead of falling back to byte-encoded row comparisons.

Yes:

- Unit tests inside \`row_backed\`: FSL / Struct roundtrip, \`take_n\`,
\`supports_type\` matches \`RowConverter::supports_fields\`.
- \`mixed_schema_column_path_uses_less_memory_than_rows_fallback\`
(mod.rs): the 54.2% memory claim + identical group assignment vs
\`GroupValuesRows\`.
- \`nested_float_edge_cases_match_rows_fallback\`: nested \`-0.0\` /
\`NaN\` produce the same groupings as \`GroupValuesRows\` (the
correctness invariant to watch, since hashing runs on the raw column and
equality runs on the row bytes).
- \`multi_batch_and_emit_first_matches_rows_fallback\`: multi-batch
streaming intern + \`EmitTo::First\` + \`take_n\`.

All 39 tests in \`aggregates::group_values\` pass.

No — internal aggregation representation only. Same query results, lower
memory footprint on mixed-schema GROUP BY keys.

- Add coverage for any type \`RowConverter\` cannot encode (currently
arrow-rs 59.x handles Map fine; \`supports_type\` delegates to
\`RowConverter::supports_fields\` so it auto-tracks upstream).
- Retire \`GroupValuesRows\` entirely once coverage is complete
(apache#23404).

(cherry picked from commit 68d5874)
<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- works towards closing apache#22682.
- replacement for apache#21765
- -
apache#21765 (review)

This PR introduces a specialized `GroupColumn` implementation for
dictionary-typed columns inside `GroupValuesColumn`, allowing dictionary
columns to participate in the columnar, vectorized aggregation path
instead of the row-based fallback.

**The Implementation is only about 175+ lines of code**. the remaining
LOC is adding extensive test at the `GroupColumn` trait level as well as
testing the `GroupValuesColumn` GroupValues trait and how it inter-opts
with multi-dictionary group by's.
<!--
Why are you proposing this change? If this is already explained clearly
in the issue then this section is not needed.
Explaining clearly why changes are proposed helps reviewers understand
your changes and offer better suggestions for fixes.
-->

- Adds a `DictionaryGroupValueBuilder` struct implementing the
`GroupColumn` trait for `Dictionary`-typed group-by columns, supporting
a configurable subset of value types
- Extends the type-check gate in `GroupValuesColumn::try_new` (the
`matches!` block) to accept `Dictionary(_, value_type)` where
`value_type` is already supported.
- Adds schema-level support so emitted dictionary group key columns
round-trip through the output schema correctly
- [removes casting
](https://github.com/apache/datafusion/blob/9e8dd76d6deb6736c51962d9c97e04be4e3f1fc9/datafusion/physical-plan/src/aggregates/group_values/multi_group_by/mod.rs#L1200)thats
done for each dictionary array in `emit`
<!--
There is no need to duplicate the description in the issue here but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->

yes. a majority of this PR is test
<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->

no. this is a pure perf boost for users.
<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.
-->

<!--
If there are any breaking changes to public APIs, please add the `api
change` label.
-->

(cherry picked from commit c1b39bd)
- Closes apache#23645
- part of apache#22715

Multi-Group-By has cases for regular Binary/LargeBinary types, but not
FixedSizeBinary

Yes.

No

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
(cherry picked from commit 3f0a953)
…ggregation-pr

cherry pick aggregation pr's

Co-authored-by: Rich-T-kid <richard.baah@datadoghq.com>
Co-authored-by: maxburke <max@urbanlogiq.com>
Co-authored-by: zhuqi-lucas <zhuqilucas@gmail.com>
## Which issue does this PR close?

Closes apache#22874.

Follow-up to [fix(topk): call attempt_early_completion when filter
rejects entire batch].

## Rationale for this change

[TopK dynamic filter pushdown attempt 2] lets `SortExec` tighten
scan-side predicates while a TopK heap finds better rows. Once the
current TopK threshold is known, scans can skip data that cannot enter
the final `ORDER BY ... LIMIT` result.

That works well for a single partition. Partitioned `SortExec` has one
extra case to handle:

- each output partition has its own local `TopK` heap
- those local heaps share one `TopKDynamicFilters` instance
- one partition can tighten the shared filter before another partition
has enough rows to fill its local heap

[fix(topk): call attempt_early_completion when filter rejects entire
batch] fixed the local case where a heap already has a max row and the
dynamic filter rejects a whole batch.

This PR fixes the remaining shared-filter case. A lagging partition can
now use the shared prefix threshold to stop early even when its local
heap is still empty. If there is no shared threshold yet, it falls back
to the existing local heap prefix check.

The shared prefix check is not treating another partition's threshold as
this partition's local heap boundary. It uses the same threshold that
already drives the shared dynamic filter. Once a partition's ordered
input has moved past that shared prefix threshold, later batches from
that partition cannot add rows that survive the shared filter. The local
heap still emits the candidates it has already kept; this only stops
pulling input that can no longer add candidates.

Single-partition behavior is unchanged.

## How the TopK optimizations fit together

There are two existing optimizations involved here:

- Dynamic filter pushdown: once a `TopK` heap has K rows, its worst kept
row becomes a threshold. That threshold tightens a scan-side filter so
later data that cannot enter the final `ORDER BY ... LIMIT` result can
be skipped.
- Prefix early exit: when the input is ordered by a prefix of the
requested sort, `TopK` can stop pulling once the last row in a batch is
past a known TopK boundary on that shared prefix.

For a partition-preserving `SortExec`, those optimizations meet in one
shared place. Each output partition has a local `TopK` heap, but all of
those local heaps publish into one shared dynamic filter. A partition
that fills first can tighten the shared filter for everyone else.
Lagging partitions then need to use that same shared prefix threshold to
stop pulling once their ordered input has moved past it.

This PR makes that composition explicit: the shared filter stays alive
until every local `TopK` has emitted, and the shared threshold carries
its common-prefix row so lagging partitions can apply the same prefix
early-exit check.

## What changes are included in this PR?

- Check early completion when a batch passes the dynamic filter but
produces zero heap replacements.
- Track local TopK emitters so a shared filter completes only after the
last emitter has produced output.
- Store the shared threshold and its common-prefix row together in
`TopKDynamicFilters`.
- Check the shared prefix in `attempt_early_completion` before falling
back to the local heap prefix.
- Add focused TopK and `SortExec` tests for the local zero-replacement
path, early completion before local heap fill, equal-prefix
non-completion, DESC/null prefix ordering, and shared-filter completion.

## How is this split for review?

The commits are ordered so each one has a narrow job:

1. `Check TopK early completion after zero-replacement batches`
Handles the local case where a batch passes the dynamic filter, produces
zero heap replacements, but still proves later rows cannot enter the
TopK.

2. `Complete shared TopK filters after all emitters`
Keeps a shared dynamic filter watchable until every local TopK emitter
has emitted. This includes the `SortExec` wiring for preserved
partitioning.

3. `Use shared TopK prefix thresholds for early exit`
Carries the shared threshold's common-prefix row so lagging partitions
can stop before their local heap is full.

## Are these changes tested?

Correctness is covered by targeted tests for:

- the local zero-replacement early-completion path
- early completion before local heap fill from a shared prefix threshold
- equal-prefix non-completion
- DESC and NULLS LAST prefix row ordering
- shared-filter completion after all TopK emitters, including preserved
`SortExec` partitioning

Relevant background:

- [perf: Add TopK benchmarks as variation over the `sort_tpch`
benchmarks] added the benchmark setup used here.
- [perf: Introduce sort prefix computation for early TopK exit
optimization on partially sorted input (10x speedup on top10 bench)]
added common-prefix TopK early termination.
- [TopK dynamic filter pushdown attempt 2] added scan-side dynamic
filter pushdown and exposed the shared-filter / local-heap interaction.
- [fix(topk): call attempt_early_completion when filter rejects entire
batch] fixed the local all-filtered-batch case.
- This PR fixes the remaining partitioned shared-filter case.

Benchmark command:

```bash
dfbench sort-tpch --sorted --limit 10 --iterations 5 \
  --path /tmp/df-topk-bench-data/tpch_sf1 \
  -o /tmp/topk-shared-prefix-followup.json
```

Both sides were rebuilt with fresh isolated `release-nonlto` target
directories before running the benchmark. This is not the regular
`topk_tpch` script: this case needs `--sorted --limit 10` to exercise
the prefix early-exit path.

Clean rerun against the PR base commit `7bb6e152b`. Times are
milliseconds, using the average of 5 iterations for each row.

| scope | PR base | this PR | change |
|---|---:|---:|---:|
| all sort-tpch queries | 760.59 | 348.84 | -54.1% |
| Q8 | 49.70 | 7.66 | -84.6% |
| Q9 | 92.13 | 10.27 | -88.8% |
| Q10 | 89.66 | 14.35 | -84.0% |

The `DataSourceExec` counters show the less noisy part of the result.
Q8/Q9/Q10 now emit only the first batch from each partition instead of
continuing to drain millions of rows.

| query | PR base `DataSourceExec output_rows` | this PR `DataSourceExec
output_rows` | PR base `bytes_scanned` | this PR `bytes_scanned` |
|---|---:|---:|---:|---:|
| Q8 | 3.66M | 81.92K | 56.81M | 15.79M |
| Q9 | 3.66M | 81.92K | 75.19M | 20.89M |
| Q10 | 3.10M | 81.92K | 110.9M | 34.69M |

## Are there any user-facing changes?

No. This is an internal physical execution optimization fix.

[perf: Add TopK benchmarks as variation over the `sort_tpch`
benchmarks]: apache#15560
[perf: Introduce sort prefix computation for early TopK exit
optimization on partially sorted input (10x speedup on top10 bench)]:
apache#15563
[TopK dynamic filter pushdown attempt 2]:
apache#15770
[fix(topk): call attempt_early_completion when filter rejects entire
batch]: apache#22852

---------

Co-authored-by: kosiew <kosiew@gmail.com>
(cherry picked from commit ba67bb4)
…rry-pick/apache-pr-22991-branch54-20260904

[branch-54] Cherry-pick apache#22991

Co-authored-by: geoffreyclaude <geoffrey.claude@datadoghq.com>
…24646)

Resolves merge conflict between HEAD's type-specific HLL accumulators and
the upstream dictionary support commit. Implements dictionary handling via
a DictionaryAccumulator wrapper that casts to the value type before
delegating to the appropriate inner accumulator.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
…T statement

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
…ict-branch

[cherry-pick] Add support for dictionary for approx_distinct (apache#24646)

Co-authored-by: Rich-T-kid <richard.baah@datadoghq.com>
Co-authored-by: mkleen <mkleen@gmail.com>
…r-defined types (apache#25015)

## Which issue does this PR close?

N/A

## Rationale for this change

`consume_user_defined_type` maps a Substrait user-defined type straight
to a bare `DataType`, so a consumer has no way to mark the resulting
`Field` for a downstream reader to recover semantic type lost in that
mapping. For example, a UDT representing JSON is commonly mapped to
plain `Utf8`, which is then indistinguishable from a real string column
once it reaches a consumer that only sees the Arrow schema (e.g. over
Arrow Flight, with no access to the original plan).

## What changes are included in this PR?

- Adds `consume_user_defined_type_metadata`, a second, additive
extension point on `SubstraitConsumer` defaulting to `Ok(None)`.
- Applies its result to the `Field` built in
`from_substrait_struct_type`, via `Field::with_metadata`.
- Existing consumers are unaffected: the default keeps today's behavior
exactly (no metadata attached).

## Are these changes tested?

Yes, two new unit tests: a consumer that supplies metadata for a
`UserDefined` field gets it attached to the field (data type unchanged),
and a consumer returning `None` (the default) attaches nothing.

## Are there any user-facing changes?

No behavior change for existing consumers. This is a new, optional trait
method with a default implementation.

(cherry picked from commit d682553)
…he#25185)

- Closes #.

`vectorized_append` is the path taken when we have non streaming hash
aggregations (when there is no ordering in the `group by`), right now it
re hashed the entire dictionary values on every batch unlike append_val
which that only does this if the values array was already seen before.

`vectorized_append` now reuses the cached value hashes when a batch
carries the same dictionary values array as the previous one (this is
the common case of a repartition or filter), instead of re hashing the
whole dictionary and re-initialising the map on every batch. This
extends the existing `Arc::ptr_eq` cache from apache#24418, which only covered
the scalar `append_val` path.

All dict tests pass, I have a PR for a benchmark that exercises this
path -> apache#25198

 no, this is a pure perf PR

(cherry picked from commit a0631ed)
@LiaCastaneda
LiaCastaneda changed the base branch from branch-54 to branch-55 September 16, 2026 10:17
@LiaCastaneda
LiaCastaneda deleted the lia.castaneda/cherry-pick/apache-pr-25185-20260916 branch September 16, 2026 10:23
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.