In Flink SQL, what does SELECT user_id, COUNT(*) FROM page_views GROUP BY user_id emit when a user's second page view arrives?
answer
- a query that never finishes
- results get corrected, not appended
- every output row carries a kind
- retract the old count first
- append-only sinks refuse updates
basics
~10 sAn update pair: -U[alice, 1] retracting the old count, then +U[alice, 2]; the first view had emitted +I[alice, 1]. Because the result updates, Flink rejects an insert-only sink for this query at planning time.
solid answer
~40 sFlink runs the query as a **continuous query** over a dynamic table, and its result is a changelog in which every row has a `RowKind`: `+I` insert, `-U` update-before, `+U` update-after, `-D` delete. The first view by `alice` creates the group and emits `+I[alice, 1]`. The second emits `-U[alice, 1]` followed by `+U[alice, 2]`, so a downstream consumer can remove the old count before applying the new one. On an insert-only source a count only grows, so no `-D` appears. The planner checks each sink's accepted changelog: a `'kafka'` table with plain `'format' = 'json'` accepts only inserts, so planning fails with `doesn't support consuming update changes which is produced by node GroupAggregate`. `EXPLAIN CHANGELOG_MODE` shows the kinds each operator produces.
code
sql · 4 linesEXPLAIN CHANGELOG_MODE
SELECT user_id, COUNT(*) AS cnt
FROM page_views
GROUP BY user_id;go deeper
Know the four row kinds and be able to write out what alice's first and second page views produce.
Explain how the planner derives changelog modes, why an insert-only sink is rejected at planning time, and how to read EXPLAIN CHANGELOG_MODE.
Trace where an updating stream starts in a real plan, pick the sink or query shape that fits the consumers, and explain the cost of retractions downstream.
Argue when consumers should receive a changelog versus final results, and how that contract should be documented for every team reading the topic.
## Dynamic tables and continuous queries Flink SQL treats a stream as a **dynamic table**: a table whose contents change as rows arrive. A query over it is a **continuous query**. It never finishes; it keeps its result table up to date as input changes. The result is itself a dynamic table, and Flink sends every change to it downstream as a **changelog**, in which each row carries a `RowKind` saying what kind of change it is. ## The four row kinds | Short form | RowKind | Meaning | |---|---|---| | `+I` | `INSERT` | a new row appears in the result | | `-U` | `UPDATE_BEFORE` | the previous version of an updated row, to be retracted | | `+U` | `UPDATE_AFTER` | the new version of an updated row | | `-D` | `DELETE` | a row disappears from the result | The short forms are how a `Row` prints its kind, for example `+I[alice, 1]`. ## Tracing the page-view count Assume `page_views` is a Kafka table with plain JSON, so its own changelog is insert-only. The aggregation keeps one accumulator per `user_id` in keyed state, and for each input row it emits the change to that user's result row: 1. `alice` views a page: the group is new, so the query emits `+I[alice, 1]`. 2. `bob` views a page: `+I[bob, 1]`. 3. `alice` views a second page: `-U[alice, 1]`, then `+U[alice, 2]`. 4. `alice` views a third page: `-U[alice, 2]`, then `+U[alice, 3]`. No `-D` appears here, because a count over an insert-only input only grows. A `-D` comes from `GROUP BY` only when the input itself retracts rows, for example an updating source, and a group loses its last row. The retraction exists so that a downstream operator can undo the old value. If a second query summed these counts, it would subtract 1 on `-U[alice, 1]` and add 2 on `+U[alice, 2]`; without the `-U` it would double-count. ## Why an insert-only sink is rejected The planner derives the **changelog mode** of every operator - the set of row kinds it can produce - and compares the final one with what the sink accepts. The Kafka connector's sink accepts whatever its value format can encode, and plain `json` encodes only inserts. Planning therefore fails before any data flows, with a `TableException` like: ```text Table sink 'default_catalog.default_database.user_counts' doesn't support consuming update changes which is produced by node GroupAggregate(groupBy=[user_id], select=[user_id, COUNT(*) AS cnt]) ``` Flink refuses rather than writing `+I[alice, 1]` and `+I[alice, 2]` as two unrelated facts that a reader could not reconcile. The same check applies when `StreamTableEnvironment.toDataStream` is called on this table: it accepts only insert-only tables. The usual ways forward are: - a sink that applies updates by key: `PRIMARY KEY (user_id) NOT ENFORCED` on an upsert-capable connector such as `'upsert-kafka'` or `'jdbc'`; - a value format that encodes a full changelog, such as `'debezium-json'` on a `'kafka'` table; - the `TO_CHANGELOG` function, new in Flink 2.3, which turns every change into an appended row with an operation-code column; - a query shape whose result is final when emitted, such as a windowed aggregation, which is a different query with different semantics. ## Seeing it yourself ```sql EXPLAIN CHANGELOG_MODE SELECT user_id, COUNT(*) AS cnt FROM page_views GROUP BY user_id; ``` The optimized physical plan attaches `changelogMode=[...]` to each node: the source scan shows `[I]`, and the `GroupAggregate` shows update kinds such as `UB` and `UA`. Reading this output is the fastest way to find which operator turns an insert-only pipeline into an updating one. ## What changes the kinds that are emitted - **The sink.** When the sink is an upsert sink keyed on `user_id`, the planner can drop `UPDATE_BEFORE` and the aggregate emits only `[I,UA]`, halving the messages. - **The input.** An updating input makes the aggregate retract accumulated values, and a group whose last row is retracted emits `-D`. - **The runtime mode.** In batch mode the same query produces one final row per user, and the changelog is insert-only. - **The query.** Adding a `HAVING cnt > 10` filter does not make the output append-only: a user crossing the threshold is inserted, and later changes still update that row.
- In Flink SQL, does adding HAVING COUNT(*) > 10 make the per-user count append-only?No. A user crossing the threshold produces `+I`, and every later view still produces `-U`/`+U` for that row. If the input could retract rows, a user falling back below the threshold would have that row retracted. `EXPLAIN CHANGELOG_MODE` confirms the plan still carries update kinds.
- What would a second Flink query that sums these counts do with the -U rows?It treats them as retractions: on `-U[alice, 1]` it subtracts 1 from its accumulator, then adds 2 on `+U[alice, 2]`. Stacked aggregations stay correct only because the retraction is carried; dropping it would double-count every update.
- How can Flink 2.3 land this updating result in an insert-only Kafka topic without a key?Wrap the query in `TO_CHANGELOG(input => TABLE ...)`. Every change becomes an appended `+I` row with an `op` column holding `INSERT`, `UPDATE_BEFORE`, `UPDATE_AFTER` or `DELETE`, and an `op_mapping` can rename or drop kinds. The consumer then owns applying those changes.
A radio scorekeeper correcting a live scoreboard says 'strike alice 1' and then 'alice 2'. A listener who can only add lines to a notebook ends up with both 1 and 2 written down as if both were true, which is why Flink refuses to hand updates to a sink that can only append.
saying these in an interview costs you the question
- Each new page view appends another row for that user
- An update is a delete followed by a fresh insert
- A kafka table with json format stores the latest count per user
- Flink requires a window before GROUP BY on a stream
- A HAVING filter turns an aggregation into append-only output