Uh oh!
There was an error while loading. Please reload this page.
fix(parquet): sync rg_plan to decoder frontier — fix wrong TopK results from re-reading already-delivered row groups (#24352) - #24354
Conversation
…pache#24352) With pushdown_filters + TopK dynamic filter pushdown, a row group whose post-predicate selection is empty is silently finished by arrow-rs without handing back a reader. rg_plan was only popped when a reader was returned, so it trailed the decoder by one; a later runtime prune then rebuilt the decoder from a stale rg_plan and re-read an already-delivered row group -- emitting duplicate rows and dropping the true top-k tail (wrong results, no error). Fix: before each prune/rebuild, sync rg_plan to the row group the decoder will actually emit next via peek_next_row_group(), dropping entries for silently finished RGs. A missing frontier RG is now an internal error rather than a silent plan drain. Adds the reporter's fixture as a regression test. Closesapache#24352.
There was a problem hiding this comment.
Pull request overview
Fixes a correctness bug in DataFusion’s Parquet push-decoder path where dynamic row-group pruning + filter pushdown could rebuild the decoder from a stale rg_plan, causing already-delivered row groups to be re-read and producing silently wrong TopK results (duplicate rows displacing the true tail).
Changes:
- Synchronize
rg_planwith the decoder’s actual next-to-emit row group (viapeek_next_row_group) before boundary prune/rebuild, preventing re-inclusion of delivered row groups. - Add an SLT regression test that reproduces the “silently-finished RG without reader” scenario and asserts correct TopK output under
parquet.pushdown_filters = true.
Reviewed changes
Copilot reviewed 2 out of 2 changed files in this pull request and generated no comments.
| File | Description |
|---|---|
| datafusion/datasource-parquet/src/push_decoder.rs | Align rg_plan to the decoder frontier at row-group boundaries to prevent stale-plan rebuilds and wrong results. |
| datafusion/sqllogictest/test_files/dynamic_row_group_pruning.slt | Adds a regression case for #24352 that fails under the buggy behavior and passes with the frontier sync fix. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@## main #24354 +/- ##
==========================================
- Coverage 81.17% 81.17% -0.01%
==========================================
Files 1109 1109 Lines 388117 388164 +47 Branches 388117 388164 +47 ==========================================
+ Hits 315071 315100 +29 - Misses 54504 54516 +12 - Partials 18542 18548 +6 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
hhhizzz
left a comment
There was a problem hiding this comment.
Thank you for your quick response! I left 2 suggestion here.
| // a reader, so without this sync `rg_plan` trails the decoder by one | ||
| // and a later rebuild can re-read an already-delivered row group | ||
| // (#24352). | ||
| if at_boundary && let Err(e) = self.sync_rg_plan_to_decoder_frontier() { |
There was a problem hiding this comment.
Could we gate this synchronization on self.row_group_pruner.is_some()?
peek_next_row_group() is not O(1): in parquet 59.2.0 it clones the remaining row-group indices and optional RowSelection, and documents a per-call cost of O(remaining row groups + selectors). With the current at_boundary-only guard, ordinary Parquet scans also call it at every row-group boundary, adding O(R²) cloning and allocation for a file with R row groups, even though no decoder rebuild can occur.
The rg_plan desynchronization is only observable when the plan is later used to rebuild the decoder, which requires row_group_pruner to be Some. Could we use:
if at_boundary
&& self.row_group_pruner.is_some()
&& letErr(e) = self.sync_rg_plan_to_decoder_frontier(){returnSome((Err(e),self));}This should preserve the #24352 fix while leaving ordinary Parquet scans unchanged. A high-row-group-count benchmark would also help validate this path.
Arrow complexity reference: https://github.com/apache/arrow-rs/blob/59.2.0/parquet/src/arrow/push_decoder/remaining.rs#L319-L334
There was a problem hiding this comment.
Good catch — done in 44f7fc7. Gated the sync on self.row_group_pruner.is_some(): only the runtime pruner rebuilds the decoder from rg_plan, so only it needs the sync, and this keeps ordinary scans (no pruner, never rebuild) from paying the peek_next_row_group() cost at every boundary.
| statement ok | ||
| RESET datafusion.explain.analyze_level; | ||
| # Regression test for #24352: TopK dynamic filter + `pushdown_filters` must not |
There was a problem hiding this comment.
Could this q26 regression explicitly enable both dynamic-filter settings and assert its own dynamic_rg_pruning=eligible plan or non-zero row_groups_pruned_dynamic_filter metric? It currently relies on defaults and final output, so a future optimizer/default change could let the test pass without exercising the prune/rebuild path that caused #24352.
There was a problem hiding this comment.
Agreed — done in 44f7fc7. The slt now enables both dynamic-filter switches explicitly instead of relying on defaults. For asserting the prune/rebuild path is actually taken, I added that to the companion Rust integration test (topk_pushdown_does_not_reread_delivered_row_group in dynamic_row_group_pruning.rs), which asserts row_groups_pruned_dynamic_filter >= 1. I went with the metric assertion there rather than a full EXPLAIN plan in the slt — the metric is a more robust guard against future optimizer/default changes and avoids a brittle plan snapshot. Happy to add an slt EXPLAIN too if you prefer belt-and-suspenders.
Adds topk_pushdown_does_not_reread_delivered_row_group in dynamic_row_group_pruning.rs (same style as the existing tests): builds the reporter's 4-row-group fixture via with_custom_data (one RG with an empty post-predicate selection invisible to statistics), and asserts the top-k has no duplicate (p4096 appears exactly once) and the true tail is present.
… tests - push_decoder: gate sync_rg_plan_to_decoder_frontier on row_group_pruner.is_some(). Only the runtime pruner rebuilds from rg_plan, and this avoids the O(remaining row groups) peek_next_row_group() cost on ordinary scans that never rebuild. - rust test: assert row_groups_pruned_dynamic_filter >= 1 so the test provably exercises the prune/rebuild path (apache#24352), not just the final output. - slt: enable both dynamic-filter switches explicitly instead of relying on defaults.
zhuqi-lucas
commented
Aug 14, 2026
Thank you @hhhizzz for review and good suggestions, addressed in latest PR, also added more test cases. |
hhhizzz
left a comment
There was a problem hiding this comment.
Thanks for the quick fix, and I hope the remaining fully-matched RowFilter optimization in can follow and land soon as well!
adriangb
commented
Aug 14, 2026
Is there any way we could structurally eliminate the possibility of drift? E.g. if we added: pubfnremaining_row_groups(&self) -> implExactSizeIterator<Item = usize> + '_Then the prune becomes decoder.remaining_row_groups().filter(|rg| !pruner.should_prune(&[*rg])), and rg_plan, RgPlanEntry, sync_rg_plan_to_decoder_frontier, advance_rg_plan_to, the pop in the Data arm, and the "we MUST build our rg_plan from this reordered list" comment in opener/mod.rs all delete. This is suggestion (2) in #24352 (comment). Or we give arrow-rs |
adriangb
left a comment
There was a problem hiding this comment.
Verified the fix rather than eyeballing it: on the merge base (0a429a37db) the new Rust test fails, and instrumenting the rebuild branch confirms it fails for the stated reason (kept=[2] pruned=1, matching @hhhizzz's Q26_REBUILD pruned=1 new_head=[2] new_len=1). On this branch it passes.
The sync is also correctly ordered and gated. peek_next_row_group() returns None whenever a row group is active, which could easily have been misread as "nothing left", and the is_at_row_group_boundary() gate in front of it rules that out. It runs on a throwaway clone of the frontier, so it does not consume the offset/limit budget.
Approving as a targeted, backportable fix. The comments below are about what the change leaves behind, not about its correctness.
Separately, and not a blocker here: the same rebuild also drops row groups out from under the carried RowSelection without slicing it, which is an independent wrong-results bug present on main with the same exposure profile. Filed as #24355 with a datafusion-cli reproduction.
| // selection is empty without handing back a reader, so without this | ||
| // sync `rg_plan` trails the decoder by one and a rebuild re-reads an | ||
| // already-delivered row group (#24352). Gating on the pruner also | ||
| // avoids the O(remaining row groups) cost of `peek_next_row_group()` |
There was a problem hiding this comment.
peek_next_row_group() clones the frontier including its Option<RowSelection>. arrow-rs documents the cost as "O(remaining row groups + selectors)". With a page-index-derived selection the selector term dominates and is paid at every boundary, so "O(remaining row groups)" understates it.
There was a problem hiding this comment.
Right — gating the sync on row_group_pruner.is_some() is exactly to avoid that peek clone on scans that never rebuild. Only the pruner rebuilds from rg_plan, so only it needs rg_plan kept in sync.
| // avoids the O(remaining row groups) cost of `peek_next_row_group()` | ||
| // on ordinary scans that never rebuild. | ||
| if at_boundary | ||
| && self.row_group_pruner.is_some() |
There was a problem hiding this comment.
This makes the correctness of rg_plan conditional on a field that has nothing to do with rg_plan's definition. Today row_group_pruner is the only consumer; the next one (a metric, a second rebuild trigger, per-RG reporting) silently reintroduces #24352. Given the perf justification above is mis-costed, I would sync unconditionally.
Bigger picture: RgPlanEntry is { rg_index: usize }, so rg_plan is a verbatim duplicate of arrow-rs's RowGroupFrontier::row_groups, and builder_from_remaining already writes that list into the rebuilt builder as row_groups: Some(row_groups) before we overwrite it with our copy. Three bugs so far have all been "the duplicate drifted": the reorder_by_statistics ordering bug, #24352, and #24355. One accessor upstream, remaining_row_groups() -> impl ExactSizeIterator<Item = usize> next to the existing row_groups_remaining() -> usize, would let rg_plan, RgPlanEntry, sync_rg_plan_to_decoder_frontier, advance_rg_plan_to and the pop in the Data arm all be deleted. That is @hhhizzz's suggested direction 2, and a smaller diff than this one. Worth a follow-up.
There was a problem hiding this comment.
Fair — the gate does couple rg_plan's correctness to row_group_pruner. It is sound today because the pruner is the only thing that rebuilds from rg_plan, but you are right that it is a fragile coupling. The clean fix is the remaining_row_groups() restructure: with no parallel rg_plan, there is no "correctness conditional on an unrelated field" to worry about.
| { | ||
| return Some((Err(e), self)); | ||
| } | ||
| if at_boundary && !self.rg_plan.is_empty() { |
There was a problem hiding this comment.
Heads up that after this fix nothing exercises the rebuild below. I put an eprintln! on the into_builder().with_row_groups(...) branch and ran all nine tests in dynamic_row_group_pruning.rs: it fires zero times. Every prune drains rg_plan and takes the return None early exit. The new test is what changes this: pre-fix the stale plan left a survivor to rebuild with, post-fix it does not.
cargo-mutants on this file agrees. Mutating pruned_count += 1 (line 365) to pruned_count *= 1 pins the count at 0, so the rebuild and the early exit never run, and all nine tests still pass, because the metric increment on the next line is a separate statement. row_groups_pruned_dynamic_filter >= 1 is therefore satisfied by the counter, not by any skipping having happened. at_boundary && ... to at_boundary || ... also survives, despite dynamic_rg_pruner_does_not_call_into_builder_mid_row_group existing for it.
Worth adding one test that prunes the middle and keeps the tail, asserting on bytes_scanned (or the arrow reader's records-read counters) rather than on the prune counter.
There was a problem hiding this comment.
Agreed — the new integration test is the first thing to actually enter that rebuild branch (pre-fix a stale survivor was left to rebuild with; post-fix not). The rebuild path's broader zero-coverage is worth its own follow-up.
| /// selection is empty (no reader handed back), which would otherwise leave | ||
| /// `rg_plan` trailing the decoder by one — a later prune/rebuild would then | ||
| /// re-include an already-delivered row group (#24352). | ||
| fn sync_rg_plan_to_decoder_frontier(&mut self) -> Result<(), DataFusionError> { |
There was a problem hiding this comment.
Result<(), DataFusionError> here vs Result<()> on advance_rg_plan_to below; datafusion_common::Result already defaults the error type.
Also: the invariant that actually broke is stated as a comment in the Data arm ("rg_plan.front() is the RG the decoder is about to read") and is still only enforced at boundaries where a pruner happens to exist. A debug_assert_eq! on that pop would have caught all three desyncs in tests that already existed, without anyone having to construct the pathological fixture, which is what @hhhizzz asked for in #24352. The cheapest upstream version is DecodeResult::Data { row_group_idx, reader }, so the caller pops to an index instead of inferring it.
There was a problem hiding this comment.
Unified to Result<()> in c4f5633. I tried the pop debug_assert_eq! too, but it false-positived: rg_plan.front() is the current RG while peek_next_row_group() returns the next one, so they only line up at the boundary where sync runs, not at every pop. Dropped it rather than keep an assert I could not time correctly — which really argues for your remaining_row_groups() idea.
| /// plan. A missing `target` means the decoder's frontier and `rg_plan` | ||
| /// have diverged; we surface that as an internal error rather than | ||
| /// silently draining the plan, which would truncate the scan. | ||
| fn advance_rg_plan_to(&mut self, target: usize) -> Result<()> { |
There was a problem hiding this comment.
Two passes over rg_plan (iter().any() then the pop loop) where one would do. Minor.
Also worth knowing this guard has no coverage: mutating e.rg_index == target to != here survives the suite.
There was a problem hiding this comment.
Folded into a single pass in c4f5633 (pop until target, error if the plan drains without finding it). And fair point that the guard itself has no coverage.
There was a problem hiding this comment.
Added coverage in bddf975: advance_rg_plan_to is now a free function over rg_plan, with two unit tests — pop-up-to-target, and the guard erroring when the target is absent. Both fail under the ==→!= mutation you flagged. And filed #24358 for the structural fix (retain_row_groups) that removes rg_plan and this whole coupling.
- sync_rg_plan_to_decoder_frontier: Result<(), DataFusionError> -> Result<()> (datafusion_common::Result already defaults the error type). - advance_rg_plan_to: fold the existence check and the pop loop into a single pass — pop until the front is `target`, and if the plan drains without finding it, return the internal error.
Thanks @adriangb, on the structural fix — agreed, |
…ction Per @adriangb, the internal_err guard on advance_rg_plan_to had no coverage (cargo-mutants: mutating == to != survived). Make advance_rg_plan_to a free function over &mut VecDeque<RgPlanEntry> so its pop/guard logic is unit-testable without a full stream state, and add two tests: pop up to target, and the guard erroring when the target is absent. Both fail under the ==/!= mutation.
zhuqi-lucas
commented
Aug 14, 2026
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
RESET reverts target_partitions to the system default (num_cpus), not the SLT runner's fixed value of 4, so the file left modified config on multi-core CI runners (4 -> 16) and failed the config-leak check. Restore it explicitly like every other slt file does.
zhuqi-lucas
commented
Aug 14, 2026
The CI is stuck, will merge it again. |
Uh oh!
There was an error while loading. Please reload this page.
…row groups (#24352) (#24368) Backport of #24354 to `branch-55` for the 55.0.0 release, per @timsaucer's request in #22393. ## Which issue does this PR close? - Backports the fix for #24352 (wrong TopK results from re-reading already-delivered row groups). ## Rationale #24352 is a **silent wrong-results** bug: with `pushdown_filters=true` + TopK dynamic filter pushdown (both on by default), a row group whose post-predicate selection is empty is finished by arrow-rs without handing back a reader, so DataFusion's `rg_plan` trails the decoder frontier by one and a later runtime prune rebuilds the decoder from a stale plan — re-reading an already-delivered row group, duplicating rows and dropping the true top-k tail. No error is raised. This is a clean cherry-pick of the squashed #24354 commit (`574fe67`); it applies to `branch-55` without conflicts. ## What changes are included? `push_decoder.rs`: sync `rg_plan` to the decoder frontier via `peek_next_row_group()` before each runtime prune/rebuild (gated on `row_group_pruner.is_some()` so ordinary scans pay nothing), with a defensive `internal_err!` if the frontier diverges from the plan. Plus the slt + rust regression tests from #24354. cc @timsaucer@alamb@adriangb
…ection is live (#24355) (#24374) Backport of #24359 to `branch-55` for the 55.0.0 release, per @timsaucer's request in #22393. Stacks cleanly on the already-merged #24368 (#24354 backport). ## Which issue does this PR close? - Backports the fix for #24355 — a second, independent silent wrong-results bug in the same parquet dynamic row-group pruning path as #24352. ## Rationale With `pushdown_filters=true` + a TopK dynamic filter, the runtime row-group pruner rebuilds the push decoder via `into_builder().with_row_groups(...)`, which drops row groups **without slicing** the carried flat page-index `RowSelection` to match — a dropped RG's selectors are then applied to the next surviving RG, silently returning wrong rows (no error). The fix declines to build the runtime `RowGroupPruner` when a row selection is present (correctness over the pruning optimization); the proper fix that keeps both is tracked upstream in apache/arrow-rs#10624 / #24358. ## Notes - Clean cherry-pick of #24359 onto `branch-55` (which now has #24354 via #24368). No conflicts. - #24359 is **approved** on `main` and pending merge; opening this now so it can ride RC3. - Verified locally on this branch: the full `dynamic_row_group_pruning` rust module (9/9) and `dynamic_row_group_pruning.slt` pass; clippy clean. cc @timsaucer@alamb@adriangb
…lection is live (apache#24355) (apache#24359) ## Which issue does this PR close? - Closesapache#24355. ## Rationale for this change With `pushdown_filters = true` + dynamic filter pushdown (on by default), a query `SELECT b FROM t WHERE <predicate on a> ORDER BY b LIMIT k` can silently return **wrong results** — rows satisfying the predicate are dropped and replaced by later ones, no error. Root cause (thanks to @adriangb's report + fixture in apache#24355): the push decoder carries one flat `RowSelection` over the concatenation of the *remaining* row groups. At a row-group boundary the runtime pruner drops row groups the dynamic predicate proves unwinnable and rebuilds the decoder: ```rust decoder.into_builder()?.with_row_groups(new_indices).build() ``` `with_row_groups(new_indices)` removes row groups **without slicing the carried `RowSelection` to match**, so the selectors intended for a dropped RG are applied to the next surviving one. In the fixture, page-index pruning leaves RG 1 with `skip 50, select 50`; after the TopK threshold prunes RG 1 and RG 2, the survivor RG 3 is decoded under RG 1's selection and its first 50 rows (`b = 0..49`, the correct answer) are wrongly skipped. This is a second, independent instance of the drift family in apache#24352/apache#24354; it is **not** fixed by apache#24354. ## What changes are included in this PR? - `opener/mod.rs`: **decline to build the runtime `RowGroupPruner` when a page-index `RowSelection` is present.** With no pruner there is no boundary rebuild, so the carried selection is never applied to the wrong row groups. This mirrors `PreparedAccessPlan::reorder_by_statistics`, which already bails when a row selection is present (`"Skipping RG reorder: row_selection present"`) because remapping the selection is too complex. This is the minimal, DataFusion-side stop-the-bleeding fix. The proper fix is upstream in arrow-rs: apache/arrow-rs#10624 proposes letting the push decoder carry **row-group-local** `RowSelection`s (`with_row_group_selections`) that are preserved across rebuilds, so dropping a row group keeps every survivor's selection aligned by construction — no global-selection slicing to get wrong, and no parallel `rg_plan` to drift (the apache#24352 path). DataFusion tracks that migration in apache#24358; this guard is removed once it lands. ## Are these changes tested? - Adds an slt regression test in `dynamic_row_group_pruning.slt` (the reporter's fixture via `generate_series` + `COPY`). It **fails on `main`** (returns `50..54` instead of `0..4`) and passes with this change. - Updates the existing rust integration test that previously asserted the runtime pruner **coexists** with a page-index selection (`dynamic_rg_pruning_coexists_with_page_index_row_selection`, `row_groups_pruned_dynamic_filter >= 1`). Since this PR intentionally disables the pruner in that case, it is renamed to `dynamic_rg_pruning_disabled_when_page_index_row_selection_present` and now asserts `row_groups_pruned_dynamic_filter == 0` while results stay correct and page-index pruning still runs. (That old test passed only because its scenario happened not to expose the bug — the misapplied selection fell outside the top-k.) - The other dynamic-prune tests are unaffected — they have no row selection, so the pruner is created as before. ## Are there any user-facing changes? Fixes silently-wrong results. Runtime row-group pruning is skipped for scans that also have a page-index row selection (correctness over a pruning optimization); this is undone once apache#24358 lands. cc @alamb@adriangb@hhhizzz
Which issue does this PR close?
Rationale for this change
With
datafusion.execution.parquet.pushdown_filters = trueand TopK dynamic filter pushdown (both on by default), a query of the shapeSELECT b FROM t WHERE <predicate on a> ORDER BY b LIMIT kcan silently return wrong results — one source row emitted several times and the true tail of the top-k missing — with no error or warning.Root cause (thanks to @hhhizzz's very detailed report + fixture in #24352): a row group whose post-predicate selection is empty is silently finished by arrow-rs without handing back a reader.
PushDecoderStreamStatepops itsrg_planonly when a reader is returned, so after a silently-finished RG the plan trails the decoder by one. When the runtime row-group pruner then rebuilds the decoder (into_builder().with_row_groups(...)) from the stalerg_plan, it re-includes an already-delivered row group, whose rows are emitted a second time and displace the genuine top-k in the heap.What changes are included in this PR?
push_decoder.rs: before each boundary prune/rebuild,rg_planis synced to the row group the decoder will actually emit next viapeek_next_row_group()(sync_rg_plan_to_decoder_frontier/advance_rg_plan_to), dropping entries for silently-finished row groups so a rebuild can never re-include a delivered group. A rebuild frontier naming an RG not in the plan is now an internal error instead of a silent plan drain.Are these changes tested?
dynamic_row_group_pruning.slt(filter columnsearch_phrasediffers from the sort columnevent_time, one row group has an empty post-predicate selection invisible to statistics). It now returns the correctp0 p4096 p4097 … p4104(was the buggyp0 p4096 p4096 …).datasource-parquetunit tests and the sqllogictest suite pass locally.Are there any user-facing changes?
Fixes silently-wrong query results; no API change.
Note
This is the standalone bug fix extracted from #23696 (per review discussion in #24352): the same
rg_plan↔ decoder-frontier sync, on its own so it merges fast and is easy to backport. #23696 will rebase on top so it carries only the fully-matchedRowFilterskip performance optimization.cc @alamb@adriangb@hhhizzz