Skip to content

Flink SQL to Kafka: fix duplicate or missing rows

Flink

A Flink SQL job that writes to Kafka usually goes wrong in one of two ways. Either the INSERT INTO statement fails before the job starts, with a planner error that names the sink and the operator upstream of it, or the job runs and the team reading the output topic reports the same order twice, a total that changes three times, or a row that never arrives. The short answer is that the changelog of the query has to match the sink and the reader: check it with EXPLAIN CHANGELOG_MODE, use the kafka connector for insert-only plans and upsert-kafka with a primary key for plans that update, and set isolation.level to read_committed on consumers of an exactly-once topic.

Table sink 'default_catalog.default_database.orders_out' doesn't support consuming update changes which is produced by node GroupAggregate(...)

A changelog is the stream of inserts, updates and deletes that a continuous Flink SQL query emits as its result table changes, and the sink connector decides how those changes are written to Kafka.

The dynamic tables documentation shows the difference with two queries. A GROUP BY count updates rows it has already emitted, so its changelog holds INSERT and UPDATE changes. A count per one-hour window only appends, so its changelog holds INSERT changes alone. Duplicate, missing and phantom rows on the Kafka side often come from a mismatch between that changelog and what the sink or the downstream reader expects.

This page is about the write side, from the Flink SQL query to the output topic. If rows are missing because the Flink job never read them, or consumer lag on the input topic looks wrong, start with Kafka consumer lag that looks wrong for a Flink job. For where Flink SQL sits in the wider system, see the complete Flink guide.

Find the symptom and jump to its section.

Read the changelog mode in the plan

Before changing the sink, find out what the query emits. EXPLAIN with the CHANGELOG_MODE detail attaches the changelog mode of every physical node to the plan, as the EXPLAIN documentation describes.

EXPLAIN CHANGELOG_MODE
INSERT INTO orders_out
SELECT customer_id, COUNT(*) AS order_count
FROM orders
GROUP BY customer_id

Read the changelogMode on the node that feeds the sink. The letters are the change kinds that the determinism documentation lists: Insert (I), Update_Before (UB), Update_After (UA) and Delete (D).

An append sink can take an insert-only plan.

changelogMode=[I]

An update plan needs an upsert sink.

changelogMode=[I,UA]

Updates with retractions and deletes need a sink that handles all four kinds.

changelogMode=[I,UB,UA,D]

Walk up the plan to the first node that adds UA or D. A regular join or a non-windowed GROUP BY is the usual one. The joins documentation notes that a regular join keeps both inputs in state forever, and that an interval join supports only append-only tables with time attributes. Rewriting the query as a window aggregation or an interval join is one way back to an insert-only result. Run EXPLAIN again after the change and confirm the sink input reads [I].

EXPLAIN also accepts PLAN_ADVICE, which the EXPLAIN documentation says is supported since Flink 1.17 and which adds warnings about risky plans. The EXPLAIN documentation shows it warning about a non-deterministic column in a pipeline that carries update messages.

EXPLAIN CHANGELOG_MODE, PLAN_ADVICE INSERT INTO ...
POV1-sqlleaky SQL on streams is a leaky abstraction

Over the years, we've sort of reflexively gone for SQL again and again. And I'm not entirely sure that's the right approach for streaming. … I think it's a very leaky abstraction.

Derek Troy-West, Co-founder and CEO of Factor House
From a podcast interview. Behind the Stream: Flink Forward Conversations (2025)

Choose kafka or upsert-kafka

The two Kafka SQL connectors take different changelogs. The Kafka SQL connector documentation lists its sink as Streaming Append Mode, so a query whose plan shows UA or D at the sink fails at planning time with the error above. The Upsert Kafka documentation lists Streaming Upsert Mode, and describes the sink writing INSERT and UPDATE_AFTER rows as normal messages and DELETE rows as messages with a null value, a tombstone for the key.

Use kafka when the plan is insert-only and every message is a new fact, such as an enriched event. Use upsert-kafka when the result is a table that changes, such as a running total per customer, and downstream readers want its latest value per key.

The dynamic tables documentation describes the two encodings behind this choice. A retract stream encodes an update as a retract message for the old row plus an add message for the new one. An upsert stream needs a unique key and encodes an update as a single message, and the consuming operator has to know the key to apply it correctly.

F1 The kafka and upsert-kafka sinks, side by side
kafka upsert-kafka
Sink mode in the docs Streaming Append Mode. Streaming Upsert Mode.
What the query may emit Inserts only. An update or delete fails at planning time. Inserts, updates and deletes.
Kafka message key Optional, set with key.fields. The PRIMARY KEY columns, which the DDL must define.
How a delete is written Not supported. A message with a null value, a tombstone for the key.
What a downstream reader must do Treat each message as a new row. Keep the last value per key and treat a null value as a delete.
Drawn from the Apache Flink stable documentation for the Kafka SQL connector and the Upsert Kafka SQL connector. Upsert Kafka connector documentation

Define the primary key and the key format

The upsert-kafka connector always works in upsert fashion and requires a primary key in the DDL. The primary key also controls which fields go into the Kafka message key, and both formats are required options.

CREATE TABLE customer_order_counts (
  customer_id STRING,
  order_count BIGINT,
  PRIMARY KEY (customer_id) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  'topic' = 'customer-order-counts',
  'properties.bootstrap.servers' = '...',
  'key.format' = 'json',
  'value.format' = 'json',
  'value.fields-include' = 'EXCEPT_KEY'
)

Three details decide whether downstream readers see the rows they expect.

The key must be the real identity of the row. Flink partitions the data on the primary key columns, so updates and deletes for the same key land in the same Kafka partition and stay in order, as the Upsert Kafka documentation describes. A primary key that does not match the GROUP BY or join key leaves several rows that a reader cannot tell apart.

The formats must match what readers deserialize. The key and value are serialized separately with key.format and value.format. With value.fields-include set to ALL, the default, the key columns are written into the value as well. With EXCEPT_KEY they appear only in the key, so a reader that looks only at the value sees rows without their key.

key.format
value.format
value.fields-include: ALL   # default

Readers must handle tombstones. A delete arrives as a message with the right key and a null value. A consumer that parses every value as JSON either fails on the null or emits an empty row, which is one source of empty rows downstream.

Why a downstream reader sees duplicates

Each cause below leaves a different pattern of keys, offsets or values on the output topic, so match it before changing a setting.

The reader treats an upsert topic as an event log

The Upsert Kafka documentation says the last record on a key takes effect when the topic is read back as a source. A consumer that counts messages instead of keeping the last value per key counts every update of a total as a new row. Reading the topic back into Flink with the upsert-kafka connector, or keeping the latest value per key in the consumer, fixes it.

To confirm it, compare the number of messages for one key with the number of distinct keys on the topic. A key with many messages and a reader whose count equals the message count is this cause.

A restart replays records under at-least-once

The delivery guarantee of the sink is set with one option, and the default for both connectors is at-least-once.

sink.delivery-guarantee: at-least-once   # default

The Kafka SQL connector documentation describes it as no records lost, although they can be duplicated. The DataStream Kafka connector documentation explains why. On a restart Flink reprocesses old input records, so everything written since the last completed checkpoint is written again. The stateful stream processing documentation describes the recovery, in which Flink selects the latest completed checkpoint and sets the sources to the stream position stored in it.

The checkpointing mode is a separate setting. With execution.checkpointing.mode set to AT_LEAST_ONCE, the same documentation says records that belong to the next checkpoint can occur as duplicates on a restore.

execution.checkpointing.mode: EXACTLY_ONCE   # default

The restart point is the checkpoint, not the consumer group offset. The same documentation says the Kafka source commits offsets when checkpoints complete and does not rely on committed offsets for fault tolerance, so the offsets visible in Kafka are for monitoring. For an upsert-kafka topic the replayed records are harmless to a reader that keeps the last value per key, which is why the Upsert Kafka documentation calls its at-least-once writes idempotent. For an append-only kafka topic every replayed record is a duplicate. A job that restarts often multiplies them, and a Flink job stuck in a restart loop covers that diagnosis.

Updates for one key arrive out of order

When the sink runs at a different parallelism from the operator above it, changelog records for a key can be shuffled out of order. The Flink table configuration reference describes two settings for sinks with primary keys. An upsert materialize operator is added before the sink by default when the planner detects disorder on unique keys, and a keyed shuffle is added by default when the sink parallelism differs from the upstream operator. Leave both at AUTO unless the plan shows a reason to change them. To confirm disorder as the cause, look for upsertMaterialize=[true] on the Sink node of the EXPLAIN output, which the EXPLAIN documentation shows on a Sink node, and compare the sink parallelism with the operator above it.

table.exec.sink.upsert-materialize: AUTO   # default
table.exec.sink.keyed-shuffle: AUTO        # default

A non-deterministic column breaks retractions

The determinism documentation explains that operators holding state process updates through complete rows, so a column whose value changes between the insert and the retraction, such as one built from CURRENT_TIMESTAMP, stops the retraction from matching. It names non-deterministic functions, a lookup join on an evolving source, and CDC metadata columns as the main causes, and describes an experimental strategy that checks the query for the problem.

table.optimizer.non-deterministic-update.strategy: TRY_RESOLVE
F2 Four causes of duplicate rows and the proof for each
What shows on the output topic What proves it
Upsert topic read as an event log Many messages for one key, and a total counted again for each update The message count for one key is higher than the number of distinct keys, and the reader keeps no last value per key
Replay under at-least-once The same record twice after a restart A restart in the job events, and duplicates written since the last completed checkpoint
Updates out of order An older value arrives after a newer one for the same key The sink parallelism differs from the operator above it, and the plan shows upsertMaterialize=[true] on the sink
Non-deterministic column A retraction that does not cancel the earlier row A column built from CURRENT_TIMESTAMP or a similar function, flagged by PLAN_ADVICE
Drawn from the Apache Flink documentation on Upsert Kafka, table configuration, EXPLAIN and determinism, all linked in this section. Determinism in continuous queries

Why rows go missing

Exactly-once delays what consumers see

With sink.delivery-guarantee set to exactly-once, the sink writes inside Kafka transactions committed on each checkpoint. The DataStream Kafka connector documentation says this delays record visibility until a checkpoint is written, for a consumer that reads only committed data. Rows written since the last checkpoint therefore look missing until it completes. The same documentation says that a consumer reading only committed data sees no duplicates after a Flink restart, which is the other half of the guarantee.

If checkpoints keep failing, the rows never appear. How to fix Flink checkpoints that fail or time out covers that case.

Backpressure and checkpoint duration are the two per-job signals to watch here. The Flink backpressure monitoring documentation rates a task HIGH when it is backpressured more than 50% of the time.

sink.delivery-guarantee: exactly-once
sink.transactional-id-prefix: orders-out-v1

The prefix is required for exactly-once, and the documentation asks for a unique prefix for each application on the same Kafka cluster so that running jobs do not interfere in each other’s transactions.

A transaction timeout shorter than the restart

The same documentation recommends a Kafka producer transaction timeout well above the maximum checkpoint duration plus the maximum restart duration, or data loss may happen when Kafka expires an uncommitted transaction. A ProducerFencedException in the job log is, in the documentation’s words, most likely a transaction timeout on the broker side.

properties.transaction.timeout.ms

State TTL drops join or aggregation results

A regular join or GROUP BY that runs forever keeps growing state, and table.exec.state.ttl clears idle state after a set time. The table configuration reference says the option sets a minimum time that idle state is retained, so with a TTL idle keys can be cleared instead of kept for as long as the job runs. The joins documentation warns that a TTL might affect the correctness of the query result, so rows that depend on cleared state can be missing from the output. How to find and fix Flink state that keeps growing covers choosing the TTL, and the default is 0, which never clears state.

table.exec.state.ttl: 0 ms   # default
F3 Three causes of missing rows and the proof for each
What shows on the output topic What proves it
Exactly-once delay Rows appear only after the next checkpoint completes sink.delivery-guarantee is exactly-once and the last completed checkpoint is older than the missing rows
Transaction timeout Records written, then lost after a long restart A ProducerFencedException in the job log, and a transaction timeout shorter than checkpoint plus restart time
State TTL Joined or aggregated rows missing after an idle period table.exec.state.ttl is above 0 and the missing rows depend on state that expired
Drawn from the Apache Flink documentation on the Kafka DataStream connector, table configuration and joins, all linked in this section. Kafka DataStream connector documentation

Phantom rows from aborted transactions

An exactly-once sink only helps consumers that read committed data. The Kafka consumer configuration reference lists read_uncommitted as the default isolation level, and says that in that mode a consumer returns all messages, even transactional messages that were aborted. A consumer on the default setting can therefore see rows from a Flink transaction that never committed. Set the isolation level on every consumer of an exactly-once topic.

The Flink documentation words the default differently. The Flink Kafka SQL connector and Upsert Kafka pages both tell readers to set isolation.level to read_uncommitted or read_committed, and describe read_committed as the default value. The Kafka reference above says read_uncommitted is the default of the Kafka consumer. A reading of the connector source on GitHub finds no code in the Kafka source or table factory that sets isolation.level, so a consumer built by the connector takes its default from the Kafka client unless the property is passed through the properties.* prefix. Rather than rely on either default, set it explicitly.

isolation.level=read_committed

For a Flink SQL table that reads the topic, pass it through the connector.

'properties.isolation.level' = 'read_committed'

The guide to Kafka producers in production covers idempotent and transactional producers from the Kafka side, and Kafka consumers in production covers the consumer side, including why duplicates still need an idempotent reader.

Where Flex fits

This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. Flex is one of several ways to watch a Flink SQL job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks.

Diagnosis. Management tools attach to Flink through its REST monitoring API, which the Flink REST API documentation says is used by Flink’s own dashboard and designed to be used by custom monitoring tools as well. Flex reads that API through the endpoint set in its Flink cluster configuration, and its system requirements say it snapshots each cluster every minute. Its job inspection documentation describes these views for each job.

  • The Topology tab shows the job’s dataflow graph, which the documentation says is identical to the Flink Job Graph, with per-subtask metrics such as Records sent.
  • The Events tab is a reverse-chronological lifecycle log with lines such as “Flink job entered state RESTARTING”, which places a restart next to the duplicates it replayed.
  • The Checkpoints tab and the Last checkpoint field show when the most recent checkpoint completed, which is when an exactly-once sink commits.
  • The jobs Overview tab counts Produced records, the total produced by sink tasks across the selected jobs.

What stays in Flink. The job documentation does not describe the changelog mode of a plan, EXPLAIN output or the options of a SQL connector, so those stay in the SQL client and the job’s DDL.

The output topic. Factor House also makes Kpow for Kafka. Its data inspect documentation describes a Key mode that searches one or more topics for an exact key.

The visible fields include Partition, Offset, Headers, Key size and Value size. Together with Key mode they list every message for one primary key and show whether a duplicate is an update, a replay or a second key.

The data produce documentation says that choosing the None serializer for the value, with a non-null key, produces a tombstone record, which is a way to test how a downstream reader handles a delete.

The inspect documentation does not say how a null value is displayed.

See restarts, checkpoints and produced records for each Flink job in Flex.

FAQ

What does “doesn’t support consuming update changes” mean in Flink SQL?

The query emits updates, usually from a regular join or a non-windowed GROUP BY, and the sink accepts only inserts. The Kafka SQL connector is an append-mode sink. Either change the query so EXPLAIN CHANGELOG_MODE shows [I] at the sink, or write to an upsert-kafka table with a primary key.

The message is built in the Flink planner’s changelog mode inference program, and the same code names delete changes as well when the query also deletes rows.

Why does upsert-kafka write duplicate messages for the same key?

At the default at-least-once guarantee, a restart replays records written since the last completed checkpoint. The Upsert Kafka documentation says the last record on a key takes effect when the topic is read back, so a reader that keeps the latest value per key is not affected. The sink buffer options keep only the last record per key between flushes, which the documentation says can also avoid some tombstone messages. Both must be set above zero to enable the buffer, and both are off by default.

sink.buffer-flush.max-rows: 0    # default, buffering off
sink.buffer-flush.interval: 0    # default, buffering off

Does exactly-once remove all duplicates from a Flink SQL Kafka sink?

Only for consumers that read committed data. The sink commits a Kafka transaction on each checkpoint, so records appear after the checkpoint, and the DataStream Kafka connector documentation says no duplicates are seen after a restart when the consumer reads only committed data. A consumer left on the Kafka consumer default of read_uncommitted can still see aborted records.

Should the Kafka message key match the Flink primary key?

For upsert-kafka it does by design, because the primary key defines the message key and the partition. For the append-only kafka connector the key comes from key.fields, and setting it to the business key keeps all messages for that key in one partition under the default partitioner. The sink partitioning section of the Kafka SQL connector documentation says the default partitioner uses a murmur2 hash of the key.

Related reading