apache / apache/datafusion

Better Grouping / aggregation pushdown

Open
#8,699 7 comments 1 reaction 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

@devinjdangelo asked on slack
https://the-asf.slack.com/archives/C04RJ0C85UZ/p1703891500973189?thread_ts=1703891392.037839&cid=C04RJ0C85UZ

> Related, is there any mechanism to push partial aggregations down to a table provider? E.g. “select count(*) from custom_provider” would push the aggregate expression down to the TableProvider rather than doing a scan.

Pushing grouping down in a plan is a classic analytic database optimization.

Among other things, it would allow doing fast data exploration like looking at counts / mins / maxes in different columns from only metadata.

### Describe the solution you'd like

I would like someone to propose
1. a change to `TableProvider` that would allow pushing aggregation down into the `TableProider`
2. A basic optimizer pass that pushes aggregation down when possible

We shouldn't innovate on API here, but should look at what other engines support in this area (e.g. TRINO / Spark / Postgres) and follow them unless there is a good reason to do something different

A good initial proof of concept would be to push `COUNT(*)`, `MIN(col)`, `MAX(col)` type queries down into parquet or other table providers that can provide such values from statistics or metadata with minimal work / data fetch

For example
```sql
datafusion-cli -c "select count(*), min(bigint_col), max(bigint_col) from './parquet-testing/data/alltypes_tiny_pages.parquet';"
DataFusion CLI v34.0.0
+----------+--------------------------------------------------------------------+--------------------------------------------------------------------+
| COUNT(*) | MIN(./parquet-testing/data/alltypes_tiny_pages.parquet.bigint_col) | MAX(./parquet-testing/data/alltypes_tiny_pages.parquet.bigint_col) |
+----------+--------------------------------------------------------------------+--------------------------------------------------------------------+
| 7300 | 0 | 90 |
+----------+--------------------------------------------------------------------+--------------------------------------------------------------------+
1 row in set. Query took 0.007 seconds.
```

The optimizer framework should eventually extend to pushing more sophisticated groupings (like pushing grouping down below joins) but many

### Describe alternatives you've considered

I believe [Eager Aggregation and Lazy Aggregation ](https://www.vldb.org/conf/1995/P345.PDF) is the classic paper on the optimizer portion

### Additional context

Contributor guide

Open the contributing guide

Research direction

Start by reviewing the TableProvider entry point and the optimizer framework, then compare aggregation pushdown support in Trino, Spark, and Postgres as requested. A first milestone would define the TableProvider change and optimizer pass, with a proof of concept for COUNT(*), MIN(col), and MAX(col) using Parquet statistics or metadata.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
data-engineering, databases, performance
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.