apache / apache/datafusion

Reorder row groups by GROUP BY keys to reduce aggregate partition state and improve cache locality

Open
#21,581 2 comments 2 reactions 0 assignees View on GitHub
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?

#21317 introduced reordering row groups by statistics for **sort pushdown** (ORDER BY / TopK queries). The same mechanism could be extended to **GROUP BY** queries to improve aggregate performance.

As @adriangb and @Dandandan suggested in https://github.com/apache/datafusion/pull/21580#discussion_r3071295010:

> Ordering by grouping keys can:
> - Reduce cardinality within partitions (partition state can be smaller)
> - Allow for better cache locality (row groups with more equal keys are grouped together)

## Describe the solution you'd like

The existing `PreparedAccessPlan::reorder_by_statistics` method already accepts any `LexOrdering` — it's generic, not tied to sort pushdown. Extending this for GROUP BY would be:

1. In the aggregate planner, compute a preferred row group ordering from the grouping keys
2. Pass it through `ParquetSource::sort_order_for_reorder`
3. Existing reorder logic handles the rest

```
SELECT category, SUM(x) FROM t GROUP BY category

Before: RGs in random order
→ HashAggregate sees mixed categories
→ Hash table stays large (all categories in memory)
→ Cache misses as different categories interleave

After: RGs ordered by category's min statistics
→ HashAggregate sees category-grouped rows
→ Can finalize/flush groups earlier (streaming aggregation opportunity)
→ Better cache locality (same category rows adjacent)
```

## Describe alternatives you've considered

- **Streaming aggregation** requires fully sorted input, which is a stronger guarantee. RG-level reorder doesn't guarantee full sort but improves locality within each RG.

## Additional context

- Parent: #21317 (row group reorder for sort pushdown)
- Related PR: #21580 (adds the `reorder_by_statistics` infrastructure)

The infrastructure from #21580 should be directly reusable — this issue is mainly about wiring up the GROUP BY planner to populate `sort_order_for_reorder` based on grouping keys.

Contributor guide

Open the contributing guide

Research direction

Start in the aggregate planner and trace how grouping keys could populate ParquetSource::sort_order_for_reorder. Read PreparedAccessPlan::reorder_by_statistics and the infrastructure from #21580, then inspect #21317 for the existing sort-pushdown path. Done means GROUP BY plans request row-group ordering by their grouping keys and the relevant aggregate behavior is covered by tests.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust, sql
Domain
data-engineering, databases, performance
Issue type
Feature
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
55/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.