A Flink 2.3 INSERT INTO fails at planning because the query's upsert key differs from the sink's primary key - what is the risk, and how do you resolve it?
answer
- two result rows, one sink key
- last writer wins, order not guaranteed
- SinkUpsertMaterializer used to hide it
- table.exec.sink.require-on-conflict
- DO ERROR, DO NOTHING, DO DEDUPLICATE
basics
~20 sSeveral result rows with different upsert keys can land on one sink key, so the stored value depends on arrival order. Flink 2.3 fails planning unless you fix the key mismatch or add ON CONFLICT DO ERROR, DO NOTHING or DO DEDUPLICATE.
solid answer
~40 sThe **upsert key** is the set of columns that uniquely identify a row of the query's result; for `GROUP BY user_id, country` it is `(user_id, country)`. If the sink declares `PRIMARY KEY (user_id) NOT ENFORCED`, the rows for `alice/DE` and `alice/FR` both write key `alice`, and whichever arrives last wins - an order that shuffling does not guarantee. Before 2.3 the planner quietly added a `SinkUpsertMaterializer` (`table.exec.sink.upsert-materialize` = `AUTO`) that kept per-key history in state. Flink 2.3 sets `table.exec.sink.require-on-conflict` to `true` by default and fails planning instead. First check for a bug: fix the query or declare `PRIMARY KEY (user_id, country)`. Otherwise choose `ON CONFLICT DO ERROR` or `DO NOTHING`, which need watermarks on every source, or `DO DEDUPLICATE`, which keeps the full history per key.
code
sql · 18 lines-- Fix the mismatch: the sink key now equals the query's upsert key
CREATE TABLE user_country_views (
user_id STRING,
country STRING,
views BIGINT,
PRIMARY KEY (user_id, country) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'user-country-views',
'properties.bootstrap.servers' = 'broker:9092',
'key.format' = 'json',
'value.format' = 'json'
);
INSERT INTO user_country_views
SELECT user_id, country, COUNT(*)
FROM page_views
GROUP BY user_id, country;go deeper
Recall that the sink's primary key should match the columns that identify each result row, usually the GROUP BY columns.
Explain what an upsert key is, how a mismatch lets two result rows overwrite one sink row, and what the three ON CONFLICT strategies do.
Diagnose the planning failure or a growing SinkUpsertMaterializer, fix the query or key first, and justify any conflict strategy by its state cost and watermark needs.
Weigh strict planning checks against upgrade friction across many pipelines, and decide which escape hatches a platform should allow and how they are reviewed.
## Upsert key versus primary key Flink's planner derives an **upsert key** for every updating result: a set of columns that identifies each result row, so that every change targets exactly one row. For a `GROUP BY`, the upsert key is the grouping columns. A sink declares its own identity with `PRIMARY KEY (...) NOT ENFORCED`. When the sink key **contains** the query's upsert key, each result row maps to one stored row and nothing extra is needed. When it does not, several result rows can map to the same stored row - a **conflict** at the sink. ## How the mismatch happens without a join ```sql CREATE TABLE user_views ( user_id STRING, country STRING, views BIGINT, PRIMARY KEY (user_id) NOT ENFORCED ) WITH ('connector' = 'upsert-kafka', ...); INSERT INTO user_views SELECT user_id, country, COUNT(*) FROM page_views GROUP BY user_id, country; ``` The query's upsert key is `(user_id, country)`; the sink's key is `user_id`. When alice views pages from two countries, the result has two live rows, `alice/DE` and `alice/FR`, and both write key `alice`. The stored value is whichever change arrives last, and changes from different upstream keys can reach the sink through different parallel paths, so that order is not deterministic. With an updating input it gets worse: retracting `alice/DE` can delete `alice` in the sink although `alice/FR` still exists. ## What Flink did before 2.3 Under `table.exec.sink.upsert-materialize` (values `NONE`, `AUTO`, `FORCE`; default `AUTO`) the planner inserted a **`SinkUpsertMaterializer`** in front of an upsert sink whenever the key mismatch could reorder changes. It kept the history of rows per primary key in state so that it could restore the right surviving row on retraction. Planning succeeded, the plan showed `upsertMaterialize=[true]`, and the state could grow without a visible warning - a classic cause of a sink operator whose state and checkpoint size kept climbing. ## What Flink 2.3 requires `table.exec.sink.require-on-conflict` defaults to `true`: if the upsert key differs from the sink's primary key and the statement has no `ON CONFLICT` clause, planning fails with a message that names both keys. The clause goes at the end of the `INSERT`: | Strategy | On a real conflict | State | Needs source watermarks | |---|---|---|---| | `ON CONFLICT DO ERROR` | fails the job at runtime | buffered changes, compacted on watermark | yes | | `ON CONFLICT DO NOTHING` | keeps the first record, discards later ones | buffered changes, compacted on watermark | yes | | `ON CONFLICT DO DEDUPLICATE` | keeps the latest, rolls back correctly on retraction | full history per key | no | `DO ERROR` and `DO NOTHING` buffer changes by primary key and upsert key and compact them when the watermark advances, so that a `-U` and a `+I` arriving out of order do not look like a conflict. Planning rejects them if any source table lacks a `WATERMARK` declaration. ## Deciding what to do 1. **Look for a logic error first.** The check exists to surface it. Here either the query should `GROUP BY user_id` only, or the sink should declare `PRIMARY KEY (user_id, country)`. 2. **If the keys are logically equal but the planner cannot prove it**, use `DO ERROR`: it costs little and fails loudly if you were wrong. 3. **If dropping later rows is acceptable**, use `DO NOTHING`. 4. **If several upstream rows genuinely update one key and correctness matters**, use `DO DEDUPLICATE`, and budget and monitor its state. ## Symptoms of the mismatch in production Before the 2.3 check, and still when the check is switched off, the mismatch shows up indirectly: - a sink value that flips between two versions for the same key from one run to the next, because the last writer depends on scheduling; - a key that disappears from the sink although the source still has live rows for it, because a retraction of one upstream row deleted the shared key; - a `SinkUpsertMaterializer` operator whose state and checkpoint size grow steadily with traffic rather than with the number of keys. Each of these points back to the same question: which columns identify a row of the result, and does the sink agree? ## Escape hatches and their cost - `table.exec.sink.require-on-conflict` = `false` restores the pre-2.3 behaviour: no clause required, and results may be non-deterministic. - `table.exec.sink.upsert-materialize` = `NONE` removes the materializer: no buffering, no compaction, no conflict resolution; records go straight to the sink in arrival order. - Neither setting fixes the underlying mismatch; both only decide who notices it and when.
- Why do ON CONFLICT DO ERROR and DO NOTHING in Flink 2.3 require watermarks on every source?Changes can reach the sink out of order, so a `-U` for one upsert key may arrive after a `+I` for another that shares the primary key, looking like a conflict. Both strategies buffer changes and compact matching pairs when the watermark advances, then judge conflicts; without watermarks they could never compact.
- What does setting table.exec.sink.upsert-materialize to NONE trade away?It removes the materializer, so there is no buffering, compaction or conflict handling; changes go to the sink in arrival order. Operators and state are lighter, but when the keys differ, the stored value per key can be wrong or non-deterministic after reordering or retraction.
- How would you spot the pre-2.3 materializer in a running Flink job?The optimized plan shows `upsertMaterialize=[true]` on the sink node, and the job graph has a `SinkUpsertMaterializer` operator in front of the sink. Its state size and checkpoint contribution growing with the number of keys is the usual symptom.
saying these in an interview costs you the question
- Upsert sinks cannot take aggregation results, hence the error
- Set upsert-materialize to NONE and the conflict is solved
- DO DEDUPLICATE is the cheapest strategy, so always pick it
- The sink primary key is enforced, so collisions cannot happen
- Key mismatches only arise from joins, never from GROUP BY