Repository navigation
[pull] main from apache:main - #280
Merged
Merged
Conversation
…atch (#22852) ## Which issue does this PR close? - Closes #22849 - A related cross-partition starvation case is tracked separately in #22874 and addressed by an upcoming follow-up PR — see [discussion](#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 #15563 (closing #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, #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>
…locations (#22918) ## Which issue does this PR close? - Closes None. ## Rationale for this change This pr simplifies heap size estimation by using a macro for types that own no heap allocations. This removes a lot of redundant code. ## What changes are included in this PR? See above. ## Are these changes tested? Yes, previous tests are passing and more tests are added. ## Are there any user-facing changes? No.
…n to the new hash aggregation impl (#22899) ## 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 #123` indicates that this PR will close issue #123. --> Part of #22710 ## 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. --> See issue for the background, this PR forward ports below optimization to the rewritten hash aggregation - #11627 After this migration, the performance is back, so this PR also changes the temporary configuration `datafusion.execution.enable_migration_aggregate` default to `true` -- the new path will be used by default. Local Clickbench_partitioned result (see `benchmarks/` for details), on M4 Pro MacBook ``` -------------------- Benchmark clickbench_partitioned.json -------------------- ┏━━━━━━━━━━━┳━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━┓ ┃ Query ┃ main ┃ split-aggr-skip-partial ┃ Change ┃ ┡━━━━━━━━━━━╇━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━┩ │ QQuery 0 │ 0.68 ms │ 0.70 ms │ no change │ │ QQuery 1 │ 7.66 ms │ 7.48 ms │ no change │ │ QQuery 2 │ 25.79 ms │ 25.72 ms │ no change │ │ QQuery 3 │ 22.25 ms │ 22.10 ms │ no change │ │ QQuery 4 │ 182.82 ms │ 188.57 ms │ no change │ │ QQuery 5 │ 213.57 ms │ 212.69 ms │ no change │ │ QQuery 6 │ 0.66 ms │ 0.69 ms │ no change │ │ QQuery 7 │ 8.54 ms │ 8.49 ms │ no change │ │ QQuery 8 │ 245.27 ms │ 246.04 ms │ no change │ │ QQuery 9 │ 323.81 ms │ 323.68 ms │ no change │ │ QQuery 10 │ 48.95 ms │ 48.70 ms │ no change │ │ QQuery 11 │ 57.73 ms │ 57.05 ms │ no change │ │ QQuery 12 │ 211.82 ms │ 210.91 ms │ no change │ │ QQuery 13 │ 298.06 ms │ 302.46 ms │ no change │ │ QQuery 14 │ 219.03 ms │ 217.94 ms │ no change │ │ QQuery 15 │ 219.24 ms │ 216.60 ms │ no change │ │ QQuery 16 │ 485.78 ms │ 493.53 ms │ no change │ │ QQuery 17 │ 500.92 ms │ 487.31 ms │ no change │ │ QQuery 18 │ 1087.29 ms │ 1051.08 ms │ no change │ │ QQuery 19 │ 19.10 ms │ 19.45 ms │ no change │ │ QQuery 20 │ 453.62 ms │ 458.61 ms │ no change │ │ QQuery 21 │ 454.90 ms │ 459.08 ms │ no change │ │ QQuery 22 │ 829.91 ms │ 847.96 ms │ no change │ │ QQuery 23 │ 2561.67 ms │ 2619.03 ms │ no change │ │ QQuery 24 │ 31.76 ms │ 31.78 ms │ no change │ │ QQuery 25 │ 86.63 ms │ 89.67 ms │ no change │ │ QQuery 26 │ 31.37 ms │ 32.67 ms │ no change │ │ QQuery 27 │ 544.97 ms │ 553.90 ms │ no change │ │ QQuery 28 │ 1822.22 ms │ 1877.44 ms │ no change │ │ QQuery 29 │ 27.76 ms │ 29.00 ms │ no change │ │ QQuery 30 │ 211.00 ms │ 217.48 ms │ no change │ │ QQuery 31 │ 206.02 ms │ 211.34 ms │ no change │ │ QQuery 32 │ 676.20 ms │ 724.32 ms │ 1.07x slower │ │ QQuery 33 │ 1144.96 ms │ 1161.21 ms │ no change │ │ QQuery 34 │ 1141.98 ms │ 1147.83 ms │ no change │ │ QQuery 35 │ 209.49 ms │ 217.06 ms │ no change │ │ QQuery 36 │ 44.38 ms │ 44.10 ms │ no change │ │ QQuery 37 │ 24.15 ms │ 24.57 ms │ no change │ │ QQuery 38 │ 29.67 ms │ 30.00 ms │ no change │ │ QQuery 39 │ 87.68 ms │ 88.80 ms │ no change │ │ QQuery 40 │ 8.57 ms │ 8.95 ms │ no change │ │ QQuery 41 │ 8.62 ms │ 8.38 ms │ no change │ │ QQuery 42 │ 7.42 ms │ 7.20 ms │ no change │ └───────────┴────────────┴─────────────────────────┴──────────────┘ ``` ## 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. --> This PR is easier to read commit-by-commit. 1. Cleanup the state machine in hash aggregation with typestate pattern 2. Move common util for partial hash aggregation skip from `aggregates/row_hash.rs` -> `aggregates/utils.rs` 3. Implement the same optimization to the migrated aggregation 4. Set configuration `enable_migration_aggregate` default to true ## 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 5. 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)? --> Existing tests + new UT ## 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. --> No
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 subscribe to this conversation on GitHub.
Already have an account?
Sign in.
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.
See Commits and Changes for more details.
Created by
pull[bot] (v2.0.0-alpha.4)
Can you help keep this open source service alive? 💖 Please sponsor : )