skip to content

In Flink SQL, how do you count ad impressions per campaign per minute with the TUMBLE window TVF, and what does the query emit?

level: juniorimportance: must knowfreq 62%

answer

  1. a table function in FROM
  2. DESCRIPTOR names the time attribute
  3. three added window columns
  4. group by window_start and window_end
  5. one final row per window

basics

~10 s

Call TUMBLE(TABLE impressions, DESCRIPTOR(impression_time), INTERVAL '1' MINUTE) in FROM and GROUP BY window_start, window_end, campaign_id. Each window emits one final insert-only row per campaign when it closes, then its state is purged.

solid answer

~40 s

A window TVF is a table-valued function called in `FROM`: `TUMBLE(TABLE impressions, DESCRIPTOR(impression_time), INTERVAL '1' MINUTE)` returns every input row plus `window_start`, `window_end` and `window_time`. Grouping by `window_start, window_end, campaign_id` is what makes the planner run a **window aggregation**. In streaming mode the descriptor must name a time attribute; with event time, each window emits one final, insert-only row per campaign once the watermark passes `window_end`, and its state is then dropped. `HOP` (slide, size), `CUMULATE` (step, size) and `SESSION` (gap, optional `PARTITION BY`) follow the same shape. The TVFs replace the legacy `GROUP BY TUMBLE(ts, ...)` group-window functions, which Flink 2.3 still documents but marks deprecated; only the TVF form supports window Top-N, window joins and window deduplication.

code

sql · 3 lines
sql
SELECT window_start, window_end, campaign_id, COUNT(*) AS impressions
FROM TUMBLE(TABLE impressions, DESCRIPTOR(impression_time), INTERVAL '1' MINUTE)
GROUP BY window_start, window_end, campaign_id;

go deeper

for a junior

Be able to write the TUMBLE query from memory: the TVF in FROM, DESCRIPTOR on the time attribute, and window_start plus window_end in GROUP BY.

for a middle

Explain why the output is insert-only and final, why window state needs no TTL, and how HOP, CUMULATE and SESSION differ in their arguments.

for a senior

Show you can pick the right TVF for a business metric, spot a query that silently fell back to a continuous GROUP BY, and migrate legacy group-window SQL.

for a principal

Weigh window TVFs against continuous aggregations for a product's freshness needs: final-per-window output is cheap and append-only, while per-row updates cost state and an upsert-capable sink.

## What a window TVF is A **window table-valued function (TVF)** takes a table and returns a table. In Flink SQL you call it in the `FROM` clause, and it assigns every input row to one or more windows by adding three columns: - `window_start` — the inclusive start of the assigned window. - `window_end` — the exclusive end of the window. - `window_time` — a time attribute equal to `window_end - 1ms`, which later time-based operations (another window TVF, an interval join) can use. The argument wrapped in `DESCRIPTOR(...)` names the column that carries time. In streaming mode it must be a **time attribute** (event time or processing time); after the TVF, that original column becomes an ordinary timestamp. A TVF on its own only labels rows. What turns it into a windowed computation is the operation on top: an aggregation, a Top-N, a join or a deduplication keyed on `window_start` and `window_end`. ## Writing the per-minute count ```sql SELECT window_start, window_end, campaign_id, COUNT(*) AS impressions FROM TUMBLE(TABLE impressions, DESCRIPTOR(impression_time), INTERVAL '1' MINUTE) GROUP BY window_start, window_end, campaign_id; ``` The planner recognises a **window aggregation** because `window_start` and `window_end` both appear in `GROUP BY`. Leave them out and it is no longer a window aggregation: grouping by `campaign_id` alone becomes an ordinary continuous `GROUP BY` that keeps a running count per campaign indefinitely and emits an update for every row. ## What the query emits A window aggregation behaves very differently from a continuous `GROUP BY`: - It emits **one final row per window and key**, not intermediate results. With an event-time attribute that happens when the watermark passes `window_end`. - Its output is **insert-only (append-only)**: a window's row is never retracted or updated, so a sink that accepts only inserts can take it. - It **purges the window's state** once the result is out, so it needs no idle-state TTL to stay bounded. - Flink's tuning guide states that mini-batch buffering is always on for window TVF aggregations, whatever `table.exec.mini-batch.enabled` says. How the watermark is generated, and what happens to rows behind it, is event-time material rather than a property of the TVF. ## The four window TVFs | TVF | Signature (Flink 2.3) | Windows per row | Typical use | |---|---|---|---| | `TUMBLE` | `TUMBLE(TABLE t, DESCRIPTOR(ts), size [, offset])` | exactly one | per-minute counts | | `HOP` | `HOP(TABLE t, DESCRIPTOR(ts), slide, size [, offset])` | several (size / slide) | trailing 10-minute totals every 5 minutes | | `CUMULATE` | `CUMULATE(TABLE t, DESCRIPTOR(ts), step, size)` | several, all sharing one start | since-midnight totals refreshed hourly | | `SESSION` | `SESSION(TABLE t [PARTITION BY k], DESCRIPTOR(ts), gap)` | one session | activity bursts per user | Two details trip people up. `HOP` takes the **slide before the size**, so `HOP(..., INTERVAL '5' MINUTES, INTERVAL '10' MINUTES)` means ten-minute windows starting every five minutes. And `SESSION` keys its sessions with `PARTITION BY` inside the call; without it, all rows share one global sequence of sessions. The session TVF is not supported in batch mode at 2.3. ## Why TVFs replaced the group-window functions Before window TVFs, Flink SQL wrote windows as grouping functions: `GROUP BY TUMBLE(impression_time, INTERVAL '1' MINUTE), campaign_id`, with `TUMBLE_START(...)` in the select list. Flink 2.3 still documents that syntax under a **deprecated** warning and recommends TVFs because: 1. The TVF form follows the SQL standard's polymorphic table functions and puts the window in `FROM`, where it composes like any other relation. 2. Only TVFs support **window Top-N**, **window joins** and **window deduplication**; group-window functions support aggregation and nothing else. 3. TVF aggregations get the optimisations in Flink's performance-tuning guide and support standard `GROUPING SETS`. ## Common mistakes - Grouping by the original time column instead of `window_start, window_end`. - Swapping `HOP`'s slide and size. - Forgetting that `window_end` is exclusive: a row stamped exactly 10:01:00 belongs to the window starting at 10:01, not to [10:00, 10:01). - Using `HOP` with a one-day size when the business wants a since-midnight running total; that is `CUMULATE`. - Pointing `DESCRIPTOR` at a plain `TIMESTAMP` column in streaming mode; it has to be a declared time attribute. - Expecting a window's row to be revised later; a window TVF aggregation emits it once.

  • What happens if you leave window_start and window_end out of the GROUP BY?
    Flink no longer plans a window aggregation. After the TVF the original time column is an ordinary timestamp, so grouping only by `campaign_id` runs a regular continuous `GROUP BY`: it keeps a count per campaign indefinitely and emits an update on every row instead of one final row per window.
  • How would you produce a since-midnight running count that refreshes every hour?
    Use `CUMULATE(TABLE impressions, DESCRIPTOR(impression_time), INTERVAL '1' HOUR, INTERVAL '1' DAY)`. Its windows are [00:00, 01:00), [00:00, 02:00) and so on up to a full day, all sharing the day's start, so each hourly row is a running total. `HOP` with a one-hour slide and one-day size would give trailing 24-hour totals instead.
  • How do you count per-user sessions that close after five quiet minutes?
    `SESSION(TABLE impressions PARTITION BY user_id, DESCRIPTOR(impression_time), INTERVAL '5' MINUTES)` with `GROUP BY user_id, window_start, window_end`. The `PARTITION BY` inside the call keys the sessions; without it every row shares one global sequence of sessions. At 2.3 the session TVF runs in streaming mode only.

saying these in an interview costs you the question

  • A window TVF aggregation emits an updated count for every incoming impression.
  • GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE) is the current recommended form in Flink 2.3.
  • Grouping by campaign_id alone after the TVF still gives per-minute counts.
  • HOP and CUMULATE are interchangeable ways to write the same sliding window.
  • A finished window's state stays until table.exec.state.ttl expires it.