skip to content

A Flink SQL GROUP BY country over an orders stream backpressures because one country dominates. How do mini-batch and two-phase aggregation help?

level: seniorimportance: should knowfreq 38%

answer

  1. one hot key, one task
  2. buffer before touching state
  3. three keys switch buffering on
  4. fold before the shuffle
  5. merge support is required

basics

~20 s

Mini-batch (table.exec.mini-batch.enabled plus allow-latency and size) buffers rows so each key's state is touched once per batch. On top of it, table.optimizer.agg-phase-strategy = 'TWO_PHASE' pre-aggregates before the shuffle, so the hot key's task receives accumulators instead of raw rows.

solid answer

~50 s

A continuous `GROUP BY` reads and writes state and emits an update for every row, and every row for the hot country lands on one task, so more parallelism does not help. **Mini-batch** — `table.exec.mini-batch.enabled = true` plus `table.exec.mini-batch.allow-latency` (say `'5 s'`) and `table.exec.mini-batch.size` (say `5000`), both required once it is on — buffers input so each key's state is accessed once per batch and fewer updates go downstream, at the cost of up to that much latency. **Two-phase (local-global) aggregation** then folds rows per key in the upstream tasks before the shuffle, so the hot task receives partial accumulators instead of millions of rows. `table.optimizer.agg-phase-strategy` defaults to `AUTO`, a cost-based choice; `TWO_PHASE` enforces it, but it depends on mini-batch and needs every aggregate, including UDAFs, to support `merge`, or Flink stays one-phase. For a skewed `COUNT(DISTINCT ...)`, also enable `table.optimizer.distinct-agg.split.enabled`.

code

sql · 8 lines
sql
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5 s';
SET 'table.exec.mini-batch.size' = '5000';
SET 'table.optimizer.agg-phase-strategy' = 'TWO_PHASE';

SELECT country, COUNT(*) AS orders, SUM(amount) AS revenue
FROM orders
GROUP BY country;

go deeper

for a junior

Know that mini-batch and two-phase aggregation exist for slow Flink SQL aggregations, and that mini-batch trades latency for throughput.

for a middle

Explain the three mini-batch keys and their defaults, how local-global aggregation folds before the shuffle, and why it depends on mini-batch.

for a senior

Diagnose a hot key from subtask metrics, enable and verify two-phase in the plan, handle distinct counts with split distinct, and anticipate the changelog side effects.

for a principal

Balance latency budgets against cost: decide which pipelines may buffer seconds of data and when a skewed metric should be redefined rather than tuned.

## Why a skewed continuous GROUP BY stalls A continuous (non-windowed) `GROUP BY` in Flink SQL is a **group aggregation**: for every input row it reads the key's accumulator from state, updates it, writes it back, and sends an update downstream. Two things go wrong when one key dominates, for example `GROUP BY country` where most orders come from one country: - **Hot task**: rows are shuffled by the grouping key, so every row for the hot country lands on one parallel task, which saturates while its peers idle, and backpressure spreads upstream. - **Per-row cost**: every row pays a state read, a state write and an emitted update; with a disk-based state backend the state access dominates. Raising parallelism does not fix the first problem, because one key always maps to one task. ## Mini-batch aggregation **Mini-batch** buffers input inside the aggregation operator and processes the buffer as a bundle, so each key's state is accessed **once per batch** instead of once per row, and one update per key per batch goes downstream. | Key | Default | Required when mini-batch is on | |---|---|---| | `table.exec.mini-batch.enabled` | `false` | `true` | | `table.exec.mini-batch.allow-latency` | `0 ms` | greater than zero, e.g. `'5 s'` | | `table.exec.mini-batch.size` | `-1` | positive, e.g. `5000` records | A batch fires when the latency interval elapses or the size is reached. The price is **added latency**, up to `allow-latency`, and burstier output. Window TVF aggregations get this buffering regardless of these keys, and Flink's tuning guide describes a mini-batch mode for regular joins too. ## Two-phase (local-global) aggregation Mini-batch reduces state access, but the hot task still receives every row. **Two-phase aggregation** splits the aggregation into a **local** phase in the upstream tasks, before the shuffle, and a **global** phase after it. Each upstream task folds its buffered rows per key into an accumulator and ships only that, so the hot country's task receives one partial result per upstream task per batch. - `table.optimizer.agg-phase-strategy` defaults to `AUTO`, which leaves the choice to the planner's cost model; `TWO_PHASE` enforces two phases and `ONE_PHASE` forbids them. - The local phase accumulates per mini-batch, so **local-global aggregation depends on mini-batch being enabled**. - Every aggregate function in the query must support merging accumulators; a user-defined `AggregateFunction` needs a `merge` method. If an aggregate call cannot be split, Flink still uses one phase. - `EXPLAIN` shows `LocalGroupAggregate` and `GlobalGroupAggregate` when it worked. ## Skewed distinct counts `COUNT(DISTINCT user_id)` grouped by a hot key gains little from local-global aggregation: a distinct accumulator must remember every distinct `user_id`, so the local accumulators are nearly as large as the raw input and the global task stays hot. **Split distinct aggregation** rewrites it into two levels, the first grouped by the original key plus a bucket computed as `HASH_CODE(user_id) % N`, the second summing the per-bucket counts: 1. Enable it with `table.optimizer.distinct-agg.split.enabled` (default `false`). 2. Size the fan-out with `table.optimizer.distinct-agg.split.bucket-num` (default `1024`). ## Side effects to check before enabling - **Latency**: every result can be delayed by up to the batch interval. - **Deduplication output**: a keep-first-row deduplication that emitted insert-only rows produces an updating changelog once mini-batch is on, which an insert-only sink downstream cannot accept. - **UDAFs without `merge`** quietly keep the aggregation one-phase. - **State size**: the global phase still holds one accumulator per key; these settings cut work and traffic, not the number of keys. ## A diagnosis order 1. Confirm the skew: one subtask of the aggregation busy and backpressured while its peers are idle. 2. Enable mini-batch with a latency the consumers can tolerate. 3. Set `TWO_PHASE` and verify the local and global operators in the plan. 4. For distinct counts, enable split distinct aggregation. 5. Re-measure the hot subtask's busy time and the end-to-end latency. ## Reading the plan Compare `EXPLAIN` output before and after changing the keys: - Without the settings, the plan shows a single `GroupAggregate` after the exchange on `country`. - With mini-batch on, a `MiniBatchAssigner` appears near the sources. - With two-phase in effect, a `LocalGroupAggregate` sits before the exchange and a `GlobalGroupAggregate` after it. If the local operator is missing, one of the preconditions failed: mini-batch is off, or an aggregate cannot merge.

  • Why does two-phase aggregation help little with COUNT(DISTINCT user_id) grouped by day, and what does?
    A distinct accumulator must keep every distinct `user_id`, so the local phase barely shrinks the data and the global task for the hot day stays overloaded. Enable `table.optimizer.distinct-agg.split.enabled`: Flink first aggregates by day plus a hash bucket of `user_id` (1024 buckets by default), then sums the bucket counts per day.
  • Why might setting TWO_PHASE have no visible effect on a query?
    Local-global aggregation depends on mini-batch, so with `table.exec.mini-batch.enabled` off there is nothing to fold locally. It also needs every aggregate to support merging; a UDAF without a `merge` method keeps the query one-phase. `EXPLAIN` confirms it: look for `LocalGroupAggregate` and `GlobalGroupAggregate`.
  • What else in the job changes when you switch mini-batch on?
    Results arrive in bursts, delayed by up to `allow-latency`. Operators that were insert-only can change too: a keep-first-row deduplication emits an updating changelog under mini-batch, so a downstream sink that accepts only inserts would be rejected by the planner.

Two-phase aggregation is like each polling station counting its own ballots and phoning in one total, instead of trucking every ballot to the central office; mini-batch is the rule that stations call in totals every few minutes rather than after each voter.

saying these in an interview costs you the question

  • Raising parallelism spreads a single hot GROUP BY key across tasks.
  • Mini-batch is on by default for every Flink SQL aggregation.
  • TWO_PHASE works on its own, without mini-batch enabled.
  • Two-phase aggregation shrinks a skewed COUNT(DISTINCT) as well as a SUM.
  • Enabling mini-batch changes latency but leaves every operator's changelog untouched.