Skip to content

fix: treat typed nulls as missing bounds in aggregate dynamic filter merge - #25149

Closed
uddhavdave wants to merge 5 commits into
apache:mainfrom
uddhavdave:fix/aggregate-dynamic-filter-typed-null-bound
Closed

uddhavdave wants to merge 5 commits into
apache:mainfrom
uddhavdave:fix/aggregate-dynamic-filter-typed-null-bound

Conversation

@uddhavdave

Copy link
Copy Markdown

Which issue does this PR close?

Rationale for this change

SELECT MIN(col), MAX(col) over a schema-evolved Parquet dataset can return a wrong MIN when datafusion.optimizer.enable_aggregate_dynamic_filter_pushdown is enabled and the query runs with multiple partitions.

When one of the files does not contain the aggregated column, the Partial aggregate for that partition evaluates to a typed null such as Int64(NULL) rather than ScalarValue::Null. The merge of per-partition bounds into the shared dynamic filter bound only short-circuited on ScalarValue::Null. The typed null therefore reached partial_cmp, where None orders before Some(_), and either replaced a valid shared minimum or blocked a later valid minimum from being recorded. The dynamic filter then lost its lower bound and became just col > <max>, which pruned the file holding the true minimum. The result depended on partition scheduling; target_partitions = 1 or disabling the pushdown returned the correct answer.

What changes are included in this PR?

  • scalar_cmp_null_short_circuit in datafusion/physical-plan/src/aggregates/aggregate_stream.rs now uses ScalarValue::is_null() so both untyped and typed nulls are treated as "no bound yet" when merging MIN/MAX bounds across partitions.

What is the testing strategy for this PR?

  • New unit test scalar_min_max_ignore_typed_nulls in aggregate_stream.rs covers scalar_min/scalar_max with typed nulls, untyped nulls, and regular values on either side. This fails without the fix and is deterministic.
  • New sqllogictest case in datafusion/sqllogictest/test_files/push_down_filter_regression.slt that builds a three-file Parquet fixture (one file without the aggregated column) and runs MIN/MAX with target_partitions = 8. The file holding the minimum is named to sort last so it is opened after the other partitions publish their bounds. With the fix reverted locally this reproduced the wrong result in 4 of 5 runs; with the fix it passes consistently. Because the failure depends on scheduling, the unit test is the primary regression guard and the slt case documents the end-to-end scenario.

cargo fmt --all, cargo clippy --all-targets --all-features -- -D warnings, and the push_down_filter_regression sqllogictest all pass locally.

Are there any user-facing changes?

No API changes. Queries affected by the bug now return the correct MIN/MAX.

@github-actions github-actions Bot added sqllogictest SQL Logic Tests (.slt) physical-plan Changes to the physical-plan crate labels Sep 10, 2026
@uddhavdave
uddhavdave force-pushed the fix/aggregate-dynamic-filter-typed-null-bound branch from b0bae23 to 0a20c96 Compare September 15, 2026 17:18

@kosiew kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@uddhavdave,

Thanks for working on this. The fix looks good overall, and the unit test gives solid deterministic coverage for the typed null merge behavior. I left one non-blocking suggestion about making the end-to-end regression more deterministic.

03)----AggregateExec: mode=Partial, gby=[], aggr=[min(agg_dyn_schema_evolution.latency_ms), max(agg_dyn_schema_evolution.latency_ms)]
04)------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_schema_evolution/01_missing.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_schema_evolution/02_high.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_schema_evolution/03_low.parquet]]}, projection=[latency_ms], file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible

query II

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The end-to-end regression here depends on partition scheduling. Repeating the query five times helps, but the old implementation could still pass if the missing-column partition does not publish its bound between the high and low partitions. The helper test gives deterministic coverage for the merge semantics, but could we add a focused integration test that controls partition completion or ordering, or otherwise forces this interleaving? That would make sure the SQL/data-source regression cannot become a false green.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think we could do this as a follow on PR as well. Or perhaps set target partitions to 1?

@uddhavdave uddhavdave Sep 24, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @kosiew for pointing this out. Agree, it's better to have a deterministic coverage. I replaced the scheduling-dependent SLT with a Rust integration test that forces the interleaving. verified that it fails deterministically on the old code.

Thank you @alamb for the suggestion. I tried setting target_partitions = 1 but had to remove the CombinePartialFinalAggregate optimizer rule to make sure aggregation mode does not get rewritten to single. The test works well.

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Sep 22, 2026
@codecov-commenter

codecov-commenter commented Sep 22, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 33.33333% with 20 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.51%. Comparing base (ed43e69) to head (6a18f43).
⚠️ Report is 2 commits behind head on main.

Files with missing lines Patch % Lines
...n/physical-plan/src/aggregates/aggregate_stream.rs 33.33% 0 Missing and 20 partials ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25149      +/-   ##
==========================================
- Coverage   82.51%   82.51%   -0.01%     
==========================================
  Files        1141     1141              
  Lines      439870   439897      +27     
  Branches   439870   439897      +27     
==========================================
  Hits       362974   362974              
- Misses      54942    54948       +6     
- Partials    21954    21975      +21     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@alamb alamb left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @uddhavdave and @kosiew

@github-actions github-actions Bot added core Core DataFusion crate and removed sqllogictest SQL Logic Tests (.slt) labels Sep 24, 2026
@uddhavdave
uddhavdave force-pushed the fix/aggregate-dynamic-filter-typed-null-bound branch from 3186f9b to 2d720ca Compare September 24, 2026 18:52
…merge

When aggregate dynamic filter pushdown is enabled with multiple
partitions, a partition whose input lacks the aggregated column (for
example a schema-evolved Parquet file) evaluates MIN/MAX to a typed null
such as `Int64(NULL)`. The shared-bound merge only short-circuited on
`ScalarValue::Null`, so the typed null fell through to `partial_cmp`,
where `None` orders before `Some(_)`, and replaced a valid shared
minimum. The dynamic filter then lost its lower bound and could prune
the file holding the true MIN, returning a wrong result depending on
partition scheduling.

Use `ScalarValue::is_null()` so both untyped and typed nulls are
ignored when merging bounds.

Closes apache#25147
@uddhavdave
uddhavdave force-pushed the fix/aggregate-dynamic-filter-typed-null-bound branch from 2d720ca to 6a18f43 Compare September 28, 2026 04:58
@github-actions github-actions Bot removed the auto detected api change Auto detected API change label Sep 28, 2026
@uddhavdave

Copy link
Copy Markdown
Author

closing this PR since it all of the code was already merged #25579

@uddhavdave uddhavdave closed this Sep 30, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core Core DataFusion crate physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

MIN dynamic filter race across partitions

4 participants