Skip to content

bench: Add drain-phase benchmark for grouped aggregation - #24795

Open
jayzhan211 wants to merge 1 commit into
apache:mainfrom
jayzhan211:agg-drain-bench
Open

bench: Add drain-phase benchmark for grouped aggregation#24795
jayzhan211 wants to merge 1 commit into
apache:mainfrom
jayzhan211:agg-drain-bench

Conversation

@jayzhan211

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

The blocked / chunked state management work in #24704 needs a way to tell
whether it worked. The suites we normally judge aggregation PRs on — ClickBench
and h2o groupby — measure whole-query time, and the phase this work changes is a
small fraction of that, so they can neither show the win nor size it. They stay
as the regression gate; this benchmark is what shows the improvement.

It measures only the drain phase: what AggregateExec does after its input is
exhausted and it starts producing output. Today both implementations drain by
calling emit(EmitTo::All) — materializing every group into one giant
RecordBatch — and then handing that batch downstream as batch.slice(..)
chunks (aggregate_hash_table/common.rs:265, labelled "temporary solution until
blocked state management is implemented", and grouped_hash_stream.rs:1337).
Two consequences, and the benchmark quantifies both:

  1. The whole drain is one poll. Time-to-first-batch equals total drain time,
    so the runtime is blocked for the entire span ([EPIC] Eliminate Long Polls in HashAggregate via Chunked Storage and Incremental Emission #19906).
  2. Nothing is released until the drain ends. Every slice shares the giant
    batch's buffers, so memory held at the halfway point is 100% of what was held
    at the start ([EPIC] Use blocked / chunked memory management in hash aggregation #24704, symptom 1).

Baseline on main (61bf6b9), 10M groups, five aggregates:

    groups       key    agg   build_ms   drain_ms   ttfb_ms   max_gap_ms  peak_pool  peak_live  pool_%  live_%
  10000000     int64   wide      284.0       16.6      16.2         16.2   1158.2MB   1168.2MB  100.0%  100.0%
  10000000      utf8   wide      512.6       16.1      15.8         15.8   1574.3MB   1616.2MB  100.0%  100.0%
  10000000      dict   wide     1035.4       33.2      32.8         32.8   1702.4MB   1798.4MB  100.0%  100.0%
  10000000   liststr   wide     1091.4     1687.1    1686.7       1686.7   2966.3MB   5226.8MB  100.0%  100.0%

Two things to read off it.

max_gap_ms ≈ ttfb_ms ≈ drain_ms in every row. The entire drain happens
inside one poll_next, exactly as #19906 describes.

Key layout decides the size of that poll, and the spread is 100x. This was
the surprise. For Int64 and Utf8 keys, emit(EmitTo::All) is close to a
buffer move — take_needed is a mem::take — so 10M groups drain in ~16 ms and
there is very little to win. Keys that go through arrow's row format have to
decode every group on the way out: GroupValuesRows::emit calls convert_rows
over the whole table (group_values/row.rs:215), and nested keys use a
row-backed GroupColumn inside the vectorized path. A List(Utf8) key at 10M
groups blocks the runtime for 1687 ms in a single poll — the ">1s stall at
~10M groups" reported in #19906, reproduced.

Worth stating explicitly for anyone measuring this work: on flat keys the
latency win is ~16 ms, so a benchmark run on Int64 keys alone will show
nothing. The four default shapes are chosen to span that range.

pool_% and live_% are 100% everywhere: nothing is released until the drain
ends, on either code path, in either accounting. Shapes not in the default set
behave the same way — Struct(Int64, Utf8) drains in 174 ms at 10M groups and
List(Int64) in 664 ms, both at 100% — and a ten-column flat key is cheap
(~2 ms at 1M groups), so column count is not what drives the stall.

What changes are included in this PR?

One new benchmark, datafusion/physical-plan/benches/aggregate_drain.rs, and
its [[bench]] entry. No changes to any non-test code.

metric meaning
drain_ms input exhausted → last output batch
ttfb_ms input exhausted → first output batch
max_gap_ms longest interval between consecutive output batches — the long-poll proxy
peak_* peak memory over the run
pool_% / live_% memory at the 50%-drained mark over memory after the first batch: ~100% means nothing is released until the drain ends, ~50% means memory is released as output is produced

Three design decisions worth flagging for review:

Criterion is not used. The quantities of interest are within-run timings and
memory samples, not a throughput distribution. The bench is harness = false
with a plain main that prints a table.

Memory is reported two ways, because they disagree. pool is what the
MemoryPool has reserved; live is bytes actually live on the heap, from a
counting global allocator. The pool does not track the materialized output
batch, so on the legacy path with utf8 keys it reports 0.1 MB reserved
after the first output batch while the process is holding 960 MB — symptom 3
in #24704. A pool-only measurement would be blind exactly where the problem is.

The ratio is anchored to the first output batch, not to the peak. Peak is
reached while the hash table is still being built, and the current code takes a
one-time step down from build state to materialized output. at50 / peak
therefore reads ~50% today and looks like incremental release already works;
at50 / first correctly reads 100%.

Input is one row per group, so the build phase is as short as possible and the
drain is what is being measured. A wrapper ExecutionPlan records when the
input is exhausted, which is what separates build from drain; a sampler thread
reads memory every 250 µs so peaks reached inside a single long poll are not
missed.

cargo bench -p datafusion-physical-plan --features test_utils --bench aggregate_drain

# just the shape that stalls
cargo bench -p datafusion-physical-plan --features test_utils --bench aggregate_drain -- \
    --keys liststr --groups 10000000 --aggs wide

# the pre-migration path
cargo bench -p datafusion-physical-plan --features test_utils --bench aggregate_drain -- --legacy

Defaults are int64, utf8, dict, liststr keys × sum, wide aggregates ×
10k, 10M groups — a cheap floor, the common case, the GroupValuesRows
fallback, and the shape that stalls. Note that liststr at 10M groups holds
several GB.

Grouped aggregation is mid-migration (#22710), and both implementations drain
the same way, so both are covered: with execution.enable_migration_aggregate
on (the default) a single grouping set runs on SingleHashAggregateStream, and
--legacy turns the flag off to measure GroupedHashAggregateStream.

How to read the output

One line per shape. The drain emits groups / batch_size batches — 1221 of them
at 10M groups with the default 8192 — so a healthy drain spreads its work across
1221 polls, and today's does not:

  • max_gap_ms is the headline. It is the longest single stretch the tokio
    worker was blocked. Compare it against the mean gap, drain_ms / batches. For
    liststr the mean gap is 1.4 ms but max_gap_ms is 1687 ms, i.e. one poll
    does 100% of the work. That ratio is the long poll.
  • ttfb_ms ≈ drain_ms is the same fact from the consumer's side. The first
    row downstream costs as much as all of them.
  • pool_% / live_% at 100% mean the operator is still holding everything
    it held at the start of the drain when it is half-finished.
  • peak_live vs peak_pool shows whether output is materialized on top of
    state rather than moved out of it. For liststr, peak_live (5227 MB) far
    exceeds peak_pool (2966 MB) because convert_rows builds the whole output
    while the row buffer is still alive; for flat keys the two nearly match
    because emit is a move.
  • build_ms is not part of the gate, but watch it: blocked storage puts a
    (block, offset) indirection on every group lookup, and that cost lands here.
    ClickBench and h2o are the real guard for it.

What the follow-up work should show

Targets for #24704 at 10M groups, derived from the baseline above. Hard gates
are marked; the rest are informational but should move in the stated direction.

metric today (liststr) expected after why
max_gap_ms 1687 < 20 (gate) one block of work per poll, not the whole table. Mean gap is 1.4 ms; anything under ~10x that is healthy
ttfb_ms 1687 < 20 (gate) first batch costs one block
live_% 100% ≤ 60% (gate) half the groups emitted ⇒ about half the memory released. 50% is ideal; the hash table itself may not shrink in step
pool_% 100% ≤ 60% (gate) must move with live_% — see failure modes below
drain_ms 1687 1400–2000 (gate: ≤ 110% of baseline) same total work, spread out. This is where #19562 died
peak_live 5227 MB ~3000–3500 MB the full-output spike disappears; only state plus one block is live. Flat keys will barely move, since they have no spike to remove

And on the flat shapes, which have almost nothing to win and everything to lose:

shape metric today requirement
utf8 @ 10M drain_ms 16.1 ≤ 18 (gate)
int64 @ 10M drain_ms 16.6 ≤ 18 (gate)
any @ 10k drain_ms ~0 no measurable regression (gate)

The strongest single check is not a threshold at all: max_gap_ms should stop
depending on group count.
Run --groups 10000000,20000000 — today max_gap_ms
roughly doubles, because the poll materializes everything. After blocked
emission it should be flat, because a poll materializes one block regardless of
how many groups exist.

Failure modes this benchmark is designed to catch:

  • drain_ms up 2x or more — EmitTo::First(n) shifting the remaining elements
    on contiguous storage, the O(remaining) trap from Incremental group emission in HashAggregate #19562.
  • live_% still 100% — blocks are being emitted but not dropped, or the output
    batches still share one allocation.
  • live_% drops but pool_% stays at 100% (or the reverse) — the reservation
    no longer describes reality. Both must move together, or downstream spill
    decisions get worse rather than better.
  • max_gap_ms down but drain_ms up — work was spread out by making more of
    it. A win on latency paid for with throughput.

Are these changes tested?

The benchmark is the test artifact; it adds no product code. cargo test --benches runs the binary with --test, which is handled as a fast smoke run
at 1k groups that asserts each configuration emits exactly one row per group.

cargo fmt --all and cargo clippy --all-targets --all-features -- -D warnings
are clean.

Are there any user-facing changes?

No. New benchmark only; no public API or behavior changes.

@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Aug 30, 2026
Adds `datafusion/physical-plan/benches/aggregate_drain.rs`, which measures
what `AggregateExec` does after its input is exhausted: total drain time,
time-to-first-batch, the longest gap between output batches (the long-poll
proxy), peak memory, and memory still held at the 50%-drained mark.

This is the measurement artifact for the blocked / chunked memory management
epic (apache#24704) and the long-poll issue (apache#19906). Both the migrated
`SingleHashAggregateStream` and the legacy `GroupedHashAggregateStream`
drain by materializing every group with `EmitTo::All` and then handing out
`batch.slice(..)` chunks, so neither releases memory until the drain ends;
`--legacy` selects the fallback path for comparison.

Key layout decides how much work the drain does. At 10M groups with five
aggregates, flat keys drain in ~16ms because `emit(EmitTo::All)` is close to
a buffer move, while a `List(Utf8)` key blocks the runtime for over a second
in a single poll: nested keys and the `GroupValuesRows` fallback have to
decode every group on the way out. The four default shapes span that range.

Memory is reported both as `MemoryPool` reservation and as live heap bytes
(via a counting global allocator), because the two disagree: the pool does
not track the materialized output batch, so it can report near-zero while
the operator still holds the whole result.

Criterion is not used - the quantities of interest are within-run timings
and memory samples rather than a throughput distribution.
@codecov-commenter

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 81.52%. Comparing base (61bf6b9) to head (28db339).
⚠️ Report is 1 commits behind head on main.

Additional details and impacted files
@@           Coverage Diff           @@
##             main   #24795   +/-   ##
=======================================
  Coverage   81.52%   81.52%           
=======================================
  Files        1123     1123           
  Lines      405970   406041   +71     
  Branches   405970   406041   +71     
=======================================
+ Hits       330978   331041   +63     
- Misses      55627    55635    +8     
  Partials    19365    19365           

☔ 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.

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

Labels

physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants