In Flink SQL, what does a streaming Top-N query for the top three products per category emit, and how do you reduce its output?
answer
- an exact ROW_NUMBER pattern
- rankings shift as sales arrive
- the unique key includes rownum
- drop rownum from the output
- the sink key must match
basics
~20 sIt emits an updating changelog: whenever the top three change, Flink retracts or updates the affected rows, keyed by category and rownum. Omitting rownum from the outer SELECT sends only the product that changed, and the sink needs a matching primary key.
solid answer
~50 sFlink recognises Top-N only in an exact pattern: `ROW_NUMBER() OVER (PARTITION BY category ORDER BY sales DESC) AS rownum` in a subquery, filtered by `rownum <= 3`. The operator keeps each category's candidates in state and, because every new order can reshuffle the ranking, its output is result-updating: changed rows go downstream as retractions or updates, never as a fixed list. With `rownum` in the output the unique key is `(category, rownum)`, so a product climbing from third to first rewrites all three rows. The **no-ranking output** optimisation omits `rownum` from the outer SELECT: then the unique key comes from upstream, here `(category, product_id)`, only the changed product's row is sent, and consumers sort the three rows themselves. Either way the sink must accept updates on the query's unique key, for example a table declared with `PRIMARY KEY (category, product_id) NOT ENFORCED`.
code
sql · 11 linesSELECT category, product_id, sales
FROM (
SELECT category, product_id, sales,
ROW_NUMBER() OVER (PARTITION BY category ORDER BY sales DESC) AS rownum
FROM (
SELECT category, product_id, SUM(amount) AS sales
FROM orders
GROUP BY category, product_id
)
)
WHERE rownum <= 3;go deeper
Know the exact Top-N pattern: ROW_NUMBER in a subquery, partitioned and ordered, filtered by rownum <= N in the outer query.
Explain why the result is updating, what its unique key is with and without rownum, and how deduplication differs by keeping the first or last row.
Diagnose a sink overwhelmed by rank rewrites, apply the no-ranking optimisation, match the sink's primary key, and consider window Top-N when the product only needs periodic results.
Challenge whether a live, continuously revised leaderboard is worth its state and write volume compared with a per-window ranking the consumers can refresh.
## How Flink recognises Top-N Flink SQL has no `TOP` keyword. It recognises a **Top-N** query from a strict pattern: `ROW_NUMBER()` over a partition in a subquery, filtered on that number. ```sql SELECT category, product_id, sales FROM ( SELECT category, product_id, sales, ROW_NUMBER() OVER (PARTITION BY category ORDER BY sales DESC) AS rownum FROM ( SELECT category, product_id, SUM(amount) AS sales FROM orders GROUP BY category, product_id ) ) WHERE rownum <= 3; ``` Rules the planner enforces: - `ROW_NUMBER()` is the only supported ranking function; `RANK()` and `DENSE_RANK()` are not supported for Top-N at 2.3. - The filter must contain `rownum <= N`; other conditions may be combined with it only by `AND`. - If the query deviates from the pattern, the optimizer cannot translate it into a Top-N. The outer query above deliberately leaves `rownum` out; the reason is below. ## Why the output is updating A batch Top-N sorts once. A streaming Top-N never finishes: every new order changes a product's `sales`, which can move it into, out of or within the top three. Flink therefore keeps each category's candidates in state and emits a **result-updating** changelog: when the top N change, the changed rows are sent downstream as retractions or updates. The input here is itself updating, because a continuous `GROUP BY` revises each product's `sales`. The Top-N operator must then process retractions of old values too. The planner picks a rank strategy from the input: a fast strategy for append-only input, and for updating input a strategy that handles retractions, which generally keeps more rows per category in state. Like any unbounded operator, its state is also subject to `table.exec.state.ttl`. ## The unique key and the no-ranking optimisation With `rownum` in the output, the unique key of the result is `(category, rownum)`. When a product climbs from third to first, the rows at ranks 1, 2 and 3 all change and all are rewritten. Across many categories and frequent reshuffles, that write amplification can make the sink the bottleneck. | Output shape | Unique key | Rows written when a product climbs from 3rd to 1st | |---|---|---| | with `rownum` | `(category, rownum)` | every shifted rank | | without `rownum` | the upstream key, here `(category, product_id)` | only the changed product's row | The **no-ranking output optimisation** is simply omitting `rownum` from the outer `SELECT`. Take category `books` with P1, P2 and P3 in the top three, and P7 overtaking P2. With `rownum` in the output, the rows keyed `(books, 2)` and `(books, 3)` are both rewritten, because different products now hold those ranks. Without `rownum`, the sink sees two changes: P7's row is added and P3's row is removed, while the rows for P1 and P2 are untouched because none of their columns changed. Consumers sort the few rows per category themselves, which is cheap. In streaming mode the sink must have the **same unique key** as the query, for example a table declared with `PRIMARY KEY (category, product_id) NOT ENFORCED` on a connector that accepts updates. A sink that accepts only inserts cannot take either shape. ## Deduplication is Top-1 on a time attribute The same pattern with `rownum = 1` and `ORDER BY` on a **time attribute** is planned as **deduplication**, a cheaper operator: 1. `ORDER BY proc_time ASC` keeps the **first** row per key. On an insert-only input with mini-batch disabled, the output is insert-only: later duplicates are simply not emitted. 2. `ORDER BY proc_time DESC` keeps the **last** row, so each newer duplicate replaces the previous one and the output is updating. 3. With mini-batch enabled, even keep-first deduplication produces an updating changelog. Ordering by a column that is not a time attribute is not deduplication; it is a general Top-1 with the heavier rank state. ## Window Top-N for per-period leaderboards If the leaderboard is per hour rather than all-time, rank the output of a window TVF aggregation with `PARTITION BY window_start, window_end, category`. **Window Top-N** emits only the final top rows at the end of each window and purges its state, so its output is insert-only and far cheaper than a never-ending updating Top-N. ## Checklist - Is the pattern exact: `ROW_NUMBER`, a subquery, `rownum <= N`? - Does anyone downstream need `rownum`, or can consumers sort? - Does the sink's primary key match the query's unique key? - Would a per-window leaderboard answer the actual question more cheaply?
- How does Flink SQL deduplication differ in output between keeping the first row and the last row?`ORDER BY proc_time ASC` with `rownum = 1` keeps the first row per key; on an insert-only input with mini-batch off it emits insert-only rows and drops later duplicates. `ORDER BY proc_time DESC` keeps the last row, so each newer duplicate replaces the previous one as an update, and the sink must accept updates.
- When would you use window Top-N instead of the continuous Top-N?When the leaderboard is per period, such as the top three per category per hour: rank the output of a window TVF aggregation with `PARTITION BY window_start, window_end, category`. It emits the final top rows once per window and purges its state, instead of an endless stream of updates.
saying these in an interview costs you the question
- Flink streaming Top-N supports RANK() and DENSE_RANK() as well as ROW_NUMBER().
- A streaming Top-N emits each top row once and never revises it.
- A sink that accepts only inserts can take a streaming Top-N result directly.
- Keeping rownum in the output makes downstream writes cheaper.
- Combining the rownum filter with OR still plans a Top-N.