apache / apache/datafusion

[EPIC] Improve window function performance (especially for large windows)

Open
#23,197 10 comments 0 reactions 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?

This ticket tracks various ideas to make evaluating window functions faster for single, often large logical window partitions. This has come up several times recently with @avantgardnerio @2010YOUY01 and @wirybeaver among others.

The core use case is to make evaluating a query like this fast:

```sql
SELECT
-- moving average over the current row and the 5 previous rows
AVG(cpu_usage) OVER (ORDER BY time ROWS BETWEEN 5 PRECEDING AND CURRENT ROW) as avg_cpu
FROM metrics;
```

The key property is that the query has one large logical window partition because there is no `PARTITION BY` clause in the window specification. Even though each frame contains only the current row and 5 preceding rows, the operator must process one global ordered window partition.

A similar problem can happen when there is a `PARTITION BY` clause but the logical window partitions are imbalanced due to skew or a small number of partition keys.

Today, each logical window partition is executed in a single DataFusion execution partition, which means:

1. It is limited to a single core
2. It is limited to a single machine in distributed environments

# Background

DataFusion has support for many window functions (see [documentation](https://datafusion.apache.org/user-guide/sql/window_functions.html)). Window functions are used like this:

```sql
SELECT
customer, time, cpu_usage,
-- moving average over the current row and the 5 previous rows
AVG(cpu_usage) OVER (PARTITION BY customer ORDER BY time ROWS BETWEEN 5 PRECEDING AND CURRENT ROW) as avg_cpu
FROM
metrics;
```

Which results in something like:

```sql
+----------+---------------------+-----------+---------+
| customer | time | cpu_usage | avg_cpu |
+----------+---------------------+-----------+---------+
| acme | 2026-01-01T00:00:00 | 10.0 | 10.0 |
| acme | 2026-01-01T00:01:00 | 20.0 | 15.0 |
...
| globex | 2026-01-01T00:03:00 | 45.0 | 30.0 |
| globex | 2026-01-01T00:04:00 | 55.0 | 35.0 |
+----------+---------------------+-----------+---------+
```

The plan for such a query looks like the following, and typically keeps all cores fully occupied:

```text
BoundedWindowAggExec(...)
SortExec(customer, time)
RepartitionExec(Hash(customer)) <-- divides work among partitions and thus cores
Scan
```

Here is an example from the tests: https://github.com/apache/datafusion/blob/a0e9887550065324320c6fd52001aa23bae67485/datafusion/sqllogictest/test_files/window_topk_pushdown.slt#L112-L118

However, for the case in question, when we remove the `PARTITION BY customer` from the `OVER` clause:

```sql
SELECT
customer, time, cpu_usage,
-- moving average over the current row and the 5 previous rows
AVG(cpu_usage) OVER (ORDER BY time ROWS BETWEEN 5 PRECEDING AND CURRENT ROW) as avg_cpu
FROM
metrics;
```

The plan looks like this (no `RepartitionExec`, and it executes in a single core):

```text
BoundedWindowAggExec(...) <-- single partition, single core
SortExec(time)
Scan
```

Here is an example: https://github.com/apache/datafusion/blob/a0e9887550065324320c6fd52001aa23bae67485/datafusion/sqllogictest/test_files/window.slt#L4100-L4103

Table definition

```sql
CREATE TABLE metrics (
time TIMESTAMP,
customer VARCHAR,
cpu_usage DOUBLE
);

INSERT INTO metrics VALUES
(TIMESTAMP '2026-01-01 00:00:00', 'acme', 10.0),
(TIMESTAMP '2026-01-01 00:01:00', 'acme', 20.0),
(TIMESTAMP '2026-01-01 00:02:00', 'acme', 30.0),
(TIMESTAMP '2026-01-01 00:03:00', 'acme', 40.0),
(TIMESTAMP '2026-01-01 00:04:00', 'acme', 50.0),
(TIMESTAMP '2026-01-01 00:00:00', 'globex', 15.0),
(TIMESTAMP '2026-01-01 00:01:00', 'globex', 25.0),
(TIMESTAMP '2026-01-01 00:02:00', 'globex', 35.0),
(TIMESTAMP '2026-01-01 00:03:00', 'globex', 45.0),
(TIMESTAMP '2026-01-01 00:04:00', 'globex', 55.0);
```

### Describe the solution you'd like

* DataFusion evaluates window functions more quickly
* DataFusion can use multiple cores to evaluate window functions for large logical window partitions
* Distributed systems such as Ballista can use multiple machines to compute window function results in parallel
* DataFusion can complete window function queries even when the logical window partitions are larger than available memory

### Describe alternatives you've considered

There are several related approaches, listed below. They are complementary rather than mutually exclusive.

### Potential Areas / follow ons

#### Improve Window Function Performance in general
- https://github.com/apache/datafusion/issues/15607 from @Dandandan
- https://github.com/apache/datafusion/issues/4904

#### Prefix Sums / Prefix Scan

Used for cumulative windows such as `ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW`.

- Paper: [Prefix Sums and Their Applications](https://www.cs.cmu.edu/~guyb/papers/Ble93.pdf)
- PoC from @avantgardnerio: https://github.com/coralogix/arrow-datafusion/pull/426

#### "Halo" Rows / Parallel Bounded Windows

Used for bounded window frames where each output partition needs some rows from neighboring ranges. Note I don't think "halo rows" is a standard database term; it means "overlapping partitions with replicated boundary rows".

- Main PoC from @avantgardnerio: https://github.com/apache/datafusion/pull/23026
- https://github.com/apache/datafusion/issues/23089
- https://github.com/apache/datafusion/pull/23090
- https://github.com/apache/datafusion/issues/23093
- https://github.com/apache/datafusion/pull/23094

#### WindowAggExec Memory / Spilling

Avoid `WindowAggExec` OOMs by processing partitions incrementally and spilling when needed.

- Support spilling for `WindowAggExec`: https://github.com/apache/datafusion/issues/22946
- PoC from @wirybeaver: https://github.com/apache/datafusion/pull/22947

#### Intra-Operator Parallelism

Another approach is to keep the logical window as one partition, but parallelize execution inside the operator.

- Paper: [Efficient Processing of Window Functions in Analytical SQL Queries](https://www.vldb.org/pvldb/vol8/p1058-leis.pdf), PVLDB 2015
- General issue from @2010YOUY01: https://github.com/apache/datafusion/issues/23174
- Window-specific issue from @2010YOUY01: https://github.com/apache/datafusion/issues/22355
- PoC from @2010YOUY01: https://github.com/apache/datafusion/pull/22356

#### Adaptive Query Execution / Runtime-Informed Planning

Relevant if DataFusion wants a general framework for runtime stats, dynamic split points, repartition choices, skew handling, and similar optimizations.

- AQE issue from @avantgardnerio: https://github.com/apache/datafusion/issues/23194
- AQE-lite PoC from @avantgardnerio: https://github.com/apache/datafusion/pull/23167

#### Related Sort / Merge Parallelism

Relevant because single-partition windows often sit downstream of sort / merge bottlenecks.

- https://github.com/apache/datafusion/pull/23124 from @Dandandan

# Related features
- Windows where the size is an expression: https://github.com/apache/datafusion/issues/15714
- Refactor the code from @2010YOUY01 : https://github.com/apache/datafusion/issues/23273
-

Contributor guide

Open the contributing guide

Research direction

Start by reading the WindowAggExec behavior and the window.slt and window_topk_pushdown.slt examples referenced in the issue. Compare the listed halo-row, spilling, intra-operator parallelism, and adaptive-query-execution approaches; done requires selecting and implementing a concrete strategy that improves large-window parallelism or memory handling.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust, sql
Domain
data, databases, distributed-systems, performance
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Quiet
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.