Skip to content

feat: per-measure FILTER predicates on Aggregate (#466) - #467

Open
Selvomega wants to merge 1 commit into
mainfrom
feat/issue-466-aggregate-filters
Open

Selvomega wants to merge 1 commit into
mainfrom
feat/issue-466-aggregate-filters

Conversation

@Selvomega

Copy link
Copy Markdown
Collaborator

Why

Closes #466.

Neither the pre- nor the post-ASAP IR could give one aggregate function its own row predicate. A predicate could only sit on the relation (Scan.predicates, Filter), where it applies to every measure. So a query that is one scan and one grouping, with one conditional count next to a plain sum, could not be planned at all:

SELECT l_shipmode, count(CASE WHEN l_returnflag = 'R' THEN 1 END), sum(l_quantity)
FROM lineitem GROUP BY l_shipmode

AggIntent::Count counts rows and never consults its argument, so lowering had no place to put the condition and rejected the query rather than silently drop it. The only alternatives were a hand-written split into two Aggregates joined back, or encoding the condition into a derived column (which only countIf does, and which erases the predicate from the IR).

What

A measure can carry a predicate with SQL FILTER (WHERE …) semantics: only rows where it is TRUE update that measure; groups are still formed from every row.

  • QueryExpr::Aggregate and ExactOperation::Aggregate gain filters: Vec<Option<Predicate>>, parallel to measures, positional against the child's output (like Filter.pred, unlike having). Empty means unfiltered, and resolve_root normalizes an all-None vector to empty so equality and CSE see one shape.
  • SummaryExpr::SummaryAgg and its executable payload gain filter: Option<Predicate>. POST_ASAP_DAG_WIRE_VERSION goes 5 → 6 so an older reader fails loudly instead of taking a filtered aggregate as unfiltered.
  • The SQL front end fills filters from three shapes: an explicit FILTER (WHERE p); count(CASE WHEN p THEN x END) → p [AND x IS NOT NULL]; and count(expr) over any other nullable expr → expr IS NOT NULL. The COUNT/corr … FILTER rejections are gone.
  • No binding rule applies a filter yet. Every recognizer in asap-aware-mapping keeps a filtered Aggregate as KeepPreAsap, and canonicalization does not promote a filtered count ranking to a heavy-hitter TopK.

Docs: docs/develop_docs/pre-asap-ir.md (the filters field, with the example above) and docs/design_docs/architecture/physical-plan-integration.md (wire version 6).

Before this PR

The query above parses, and lowering rejects it:

unsupported feature: COUNT of a nullable expression or with FILTER requires explicit per-aggregate null/filter semantics

FILTER (WHERE …) did not even reach the planner: sqlparser's GenericDialect has supports_filter_during_aggregation off, and DataFusion only selects a dialect by name.

After this PR

The same query lowers to one Aggregate over one Scan, no Join, no derived column. The CASE condition became the Count's filter; the Sum is unfiltered:

Project([0, 1, 2])
  Aggregate(
      reduction = Reduce(by = [0]),                       -- l_shipmode
      measures  = [Count, Sum(col: 2)],                   -- l_quantity
      filters   = [Some(Column(1) == 'R'), None],         -- l_returnflag
      child     = Scan("lineitem"),
  )

keep_pre_asap on that tree yields an exact plan and compiles to an executable DAG. SELECT sum(bytes) FILTER (WHERE service = 'a'), count(*) … lowers the same way, with filters = [Some(service == 'a'), None].

How

Crate Change
asap-types filters on QueryExpr::Aggregate / ExactOperation::Aggregate; filter on SummaryAgg and ExecutableOperatorPayload::SummaryAgg; wire version 6. resolve.rs binds each filter against the child schema and normalizes all-None to empty. cse.rs (structural hash and rebuild), schema_resolver.rs (usage-derived catalog), dag_export.rs and the post-ASAP CSE identity carry the field. canonicalize.rs skips filtered rankings. execution_data_state.rs counts filter columns as referenced when checking exact-operator inputs. any_measure_filtered() is the shared gate.
asap-frontend-sql New sql/dialect.rs: GenericDialect with only supports_filter_during_aggregation flipped; lower parses through DFParser::parse_sql_with_dialect + statement_to_plan, as the ClickHouse path already did. lower_aggregate computes one measure_filter per aggregate call, passes the filter's columns through any derived-column Project, and emits filters. Temporal aggregates and ROLLUP/GROUPING SETS cannot carry a filter and now reject one explicitly instead of dropping it.
asap-aware-mapping bindable_intent, exact composition, rollup, accuracy reconciliation, maintained population, the avg and composed rewrites, exact top-k, query_time_nested_sum and the analytical cost lowering all decline a filtered Aggregate. SummaryAgg relinking carries filter.
sql-function-catalog countIf still lowers to sum(CASE …); its doc now says moving the -If family onto filters is a follow-up.

Design choices, as in the issue: a parallel vector rather than a Measure struct, so existing readers of measures keep compiling; a field on SummaryAgg rather than a Filter child, so summaries that differ only in predicate can still share one child.

Tests

  • Serde: filters round-trips, and an Aggregate serialized before the field existed still reads as unfiltered (aggregate_filters_serde_round_trip_and_default).
  • CSE: two aggregates that differ only in one predicate are not shared (filtered_and_unfiltered_aggregates_do_not_merge).
  • Lowering: the example above lowers to one Aggregate with no Join (conditional_count_lowers_to_a_filtered_measure); FILTER (WHERE …) lands on the measure it annotates; count(nullable) becomes IS NOT NULL; the filter's columns survive a derived-column Project; a filter inside ROLLUP is rejected; corr … FILTER now lowers.
  • Resolution: filters bind positionally against the child; all-None normalizes to empty.
  • Canonicalization: a filtered count ranking stays Sort + Limit.
  • Mapping: a filtered measure is retained as KeepPreAsap (filtered_measure_stays_logical).
  • Dialect: mirrors every GenericDialect switch and parses the FILTER clause the generic dialect rejects.
  • Two existing tests changed with the semantics: corr_rejects_unrepresented_forms no longer lists FILTER; count_preserves_non_null_inputs_and_rejects_erased_null_semantics is now count_null_semantics_become_a_measure_filter.

cargo test --workspace, cargo clippy --workspace --tests and cargo fmt are clean.

Not in this PR

  • Binding filtered measures to summaries (SummaryAgg.filter is never set by a rule yet).
  • NULL semantics of unfiltered inputs: sum(CASE WHEN p THEN x END) still lowers to a derived column and relies on undefined NULL-skipping.
  • ClickHouse countIf/sumIf/avgIf/… on filters.

🤖 Generated with Claude Code

Neither IR could give one aggregate function its own row predicate, so
`count(CASE WHEN p THEN 1 END)` next to a plain `sum(x)` was rejected at
lowering. `QueryExpr::Aggregate` and `ExactOperation::Aggregate` gain
`filters: Vec<Option<Predicate>>`, parallel to `measures` and positional
against the child; `SummaryAgg` and its wire payload gain `filter`, and
the executable DAG wire version goes to 6.

The SQL front end fills the field from an explicit `FILTER (WHERE …)`
(parsed through a GenericDialect wrapper that enables the clause), from
`count(CASE WHEN p THEN x END)`, and from `count(expr)` over a nullable
`expr`. Resolution, canonicalization, CSE, dependency collection and DAG
export carry it. No binding rule applies a filter yet: every recognizer in
asap-aware-mapping keeps a filtered Aggregate as `KeepPreAsap`, and the
heavy-hitter promotion skips filtered rankings.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

@milindsrivastava1997 milindsrivastava1997 left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@Selvomega Thanks this is useful. Few questions:

  • What is the behavior when multiple aggregations share the same filter? Is there any logic to deduplicate filters?
  • What is the IR when there is only a single aggregate with a filter statement and what is the IR when there is only a single aggregate with an equivalent case statement? Do these equivalent SQL queries produce the same or different IRs?

If the above case produces different IRs, we may add a canonicalization pass that normalizes them to the same IR. I believe you will find such existing canonicalization logic in the codebase.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Design] Supporting Predicates in Aggregation Functions

2 participants