skip to content

In Flink SQL, how do you write a continuously updated per-user count to Kafka or JDBC, and what does PRIMARY KEY NOT ENFORCED change?

level: middleimportance: should knowfreq 54%

answer

  1. the sink needs to know identity
  2. a key Flink never checks
  3. overwrite by key, no retraction needed
  4. upsert-kafka needs key and value formats
  5. deletes become null-value records

basics

~20 s

Declare the sink with PRIMARY KEY (user_id) NOT ENFORCED on an upsert-capable connector such as 'upsert-kafka' or 'jdbc'. The key tells Flink which row each change targets, so updates overwrite by key; NOT ENFORCED means Flink never verifies uniqueness.

solid answer

~40 s

A continuous `GROUP BY user_id` emits updates, so the sink must be able to apply them. Declaring `PRIMARY KEY (user_id) NOT ENFORCED` gives the sink an identity for each row. `NOT ENFORCED` is the only mode Flink supports: it does not own the data and never checks uniqueness, so the query must guarantee it. With `'connector' = 'upsert-kafka'` (which requires the key plus `'key.format'` and `'value.format'`), `+I` and `+U` become records keyed by `user_id`, and `-D` becomes a record with a null value, a tombstone. With `'connector' = 'jdbc'`, a declared key switches the sink from append mode to upsert mode. Because an upsert sink overwrites by key, the planner can drop `UPDATE_BEFORE` rows entirely and send only `[I,UA]` - half the messages of a retract stream.

code

sql · 16 lines
sql
CREATE TABLE user_view_counts (
  user_id STRING,
  views   BIGINT,
  PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  'topic' = 'user-view-counts',
  'properties.bootstrap.servers' = 'broker:9092',
  'key.format' = 'json',
  'value.format' = 'json'
);

INSERT INTO user_view_counts
SELECT user_id, COUNT(*) AS views
FROM page_views
GROUP BY user_id;

go deeper

for a junior

Remember that an updating query needs a sink with PRIMARY KEY ... NOT ENFORCED and an upsert-capable connector such as upsert-kafka or jdbc.

for a middle

Explain retract versus upsert encoding, why the planner can drop UPDATE_BEFORE for an upsert sink, and how upsert-kafka writes deletes as tombstones.

for a senior

Match the sink key to the query's upsert key, anticipate ChangelogNormalize state when the topic is read back, and make sure consumers understand tombstones.

for a principal

Decide whether a shared topic should publish upserts, a full changelog or append-only change records, weighing consumer simplicity against topic size and state cost.

## The problem an upsert sink solves A continuous `SELECT user_id, COUNT(*) FROM page_views GROUP BY user_id` in Flink SQL emits a changelog: `+I[alice, 1]`, then `-U[alice, 1]` and `+U[alice, 2]`, and so on. A sink that can only append would store every version as a separate fact, so Flink rejects it at planning time. The sink needs to know **which stored row each change targets**, and that is what a primary key declares. ## Declaring the key: PRIMARY KEY ... NOT ENFORCED ```sql CREATE TABLE user_view_counts ( user_id STRING, views BIGINT, PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( ... ); ``` - The constraint says the listed columns are unique and not null. Declaring it also makes those columns `NOT NULL`. - **`NOT ENFORCED` is the only supported mode.** Flink does not own the data in Kafka or the database, so it never validates uniqueness; the query is responsible for producing at most one live row per key. - The planner uses the key: it lets an upsert-capable connector apply changes by key, and it compares the key with the **upsert key** the query derives (here `user_id`, the grouping key). ## Retract versus upsert | | Retract stream | Upsert stream | |---|---|---| | Messages per update | two (`-U`, `+U`) | one (`+U`) | | Needs a unique key | no | yes | | Carries `UPDATE_BEFORE` | yes | no | | Consumer applies | add and remove rows | overwrite or delete by key | An upsert consumer that receives `+U[alice, 2]` simply overwrites the row for `alice`, so the retraction carries no information it needs. When the sink is an upsert sink and its key matches the query's upsert key, the planner removes `UPDATE_BEFORE` from the whole path; `EXPLAIN CHANGELOG_MODE` then shows the aggregate producing `[I,UA]`. ## Connector by connector - **`'connector' = 'upsert-kafka'`** requires `PRIMARY KEY`, `'key.format'` and `'value.format'`. The key columns become the Kafka record key. `INSERT` and `UPDATE_AFTER` are written as normal records; `DELETE` is written as a record with a null value, a tombstone for that key. Flink partitions by the key columns, so all changes to one key land in one partition in order. - **`'connector' = 'jdbc'`** runs in upsert mode when a key is declared and in append mode otherwise. Given an updating query and no key, planning fails with `please declare primary key for sink table when query contains update/delete record.` The declared key should match a primary or unique key of the database table. - **`'connector' = 'kafka'` with plain `'format' = 'json'`** accepts only inserts, so it is rejected for this query. - **`'connector' = 'kafka'` with `'format' = 'debezium-json'`** accepts a full changelog; Flink writes `UPDATE_BEFORE` and `UPDATE_AFTER` as a delete and an insert message. - **`TO_CHANGELOG`**, new in Flink 2.3, turns an updating table into append-only rows with an operation-code column, for sinks that can only append. ## What happens on the read side Reading an `upsert-kafka` topic back into Flink produces only `UPDATE_AFTER` and `DELETE` rows: a record is an upsert, and a null value is a delete. If a downstream operator needs the full changelog with `UPDATE_BEFORE`, the planner inserts a stateful **`ChangelogNormalize`** operator that remembers the last value per key to produce the missing retractions. Its state grows with the number of distinct keys - the cost of the compact topic format. ## Checklist 1. Confirm the query's result has a unique key: here the `GROUP BY` columns. 2. Declare that same key as `PRIMARY KEY (...) NOT ENFORCED` on the sink. 3. Choose an upsert-capable connector, and for `upsert-kafka` set both `key.format` and `value.format`. 4. Check the plan with `EXPLAIN CHANGELOG_MODE`: the path into the sink should show `[I,UA]`, with `D` only if the input can retract. 5. Make sure the consumers of the sink understand tombstones and last-value-per-key semantics. ## Mistakes that show up in review - Declaring a key on the sink that is not the query's grouping key, for example `PRIMARY KEY (page_url)` on a per-user count: several result rows then collide on one sink row, and Flink 2.3 refuses to plan the statement without an `ON CONFLICT` clause. - Reading `NOT ENFORCED` as permission to skip the key: without a key, `'jdbc'` falls back to append mode and `'upsert-kafka'` refuses to create the table. - Letting consumers treat the upsert topic as an event log: a reader that counts records instead of keeping the last value per key will count every intermediate total. - Forgetting that deletes are null-value records: a consumer that fails on a null value breaks the first time the input retracts a group.

  • What happens if the same Flink query writes to a 'jdbc' sink that declares no primary key?
    Planning fails with `please declare primary key for sink table when query contains update/delete record.` Without a key the JDBC sink runs in append mode and treats every row as an insert, which cannot express updates. Declare the key, and match it to a unique key of the database table.
  • Another Flink job reads the upsert-kafka topic and feeds a retract-capable consumer. What does the planner add, and what does it cost?
    The topic yields only `UPDATE_AFTER` and `DELETE`, so the planner adds a `ChangelogNormalize` operator that keeps the last value per key in state to produce `UPDATE_BEFORE`. The state grows with the number of distinct keys, which is the price of the compact upsert format.
  • Why must the sink's primary key match the GROUP BY columns?
    The grouping columns are the query's upsert key: they identify each result row. If the sink key matches them, every change maps to exactly one stored row. If it does not, several result rows can collide on one sink key, and Flink 2.3 then demands an explicit `ON CONFLICT` strategy.

saying these in an interview costs you the question

  • NOT ENFORCED means Flink ignores the primary key
  • Flink checks that incoming rows never repeat a key
  • Any 'kafka' table with json can store a GROUP BY result
  • upsert-kafka writes deletes as records with a delete flag
  • An upsert sink still needs every -U row to stay correct