Repository navigation
Cherry-pick apache/datafusion#25185 - #182
Closed
LiaCastaneda wants to merge 81 commits into
Closed
LiaCastaneda wants to merge 81 commits into
LiaCastaneda wants to merge 81 commits into
Conversation
…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 434957e - turn off submodule updating
Cherry pick a0763db - Fix ArrayCompact incompatibility
Cherry pick 6692f6f - fix(substrait): dedupe names
Cherry pick 12d6c81 - Add lambda substrait support (apache#21193) (apache#134)
…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
deleted the
lia.castaneda/cherry-pick/apache-pr-25185-20260916
branch
September 16, 2026 10:23
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
cherry-picks apache#25185
Also fixes a cache-invalidation gap in
take_nexposed by conflict resolution with DataDog's pre-existing partial port of this optimization (#171):cached_valueswas not cleared aftertake_nremaps inner slots, so a stale cache could be reused against remapped data. Addedself.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.