skip to content

Why does end-to-end exactly-once in Flink require a two-phase-commit sink?

level: seniorimportance: must knowfreq 62%

answer

  1. the engine cannot un-send a write
  2. two hooks, snapshot and completion
  3. stage now, publish only when confirmed
  4. what if it dies between the two
  5. the reader must ignore uncommitted data

basics

~20 s

Checkpoints only roll back Flink's internal state; output already sent to an external system cannot be un-sent when the job replays. A two-phase-commit sink stages output in a transaction, pre-commits it at the checkpoint, and commits only once Flink confirms that checkpoint completed.

solid answer

~50 s

Flink's checkpoint guarantee is about **state**, not delivery. On recovery the sources rewind and records are reprocessed, so a sink that wrote them straight through has already published duplicates that no rollback can retract. The fix is to tie the external commit to the checkpoint lifecycle: the sink writes records into an open transaction, **pre-commits** (flushes and makes durable but not visible) when the barrier triggers its snapshot, and **commits** only after Flink reports the checkpoint complete. If the job dies between the two, recovery restores the pending transaction from the checkpoint and commits it; if it dies before pre-commit, the transaction is aborted and the replay reproduces its records. In Flink 2.3 this is the Sink V2 contract — a `CommittingSinkWriter` hands committables to a `Committer` — which replaced the `SinkFunction`-based `TwoPhaseCommitSinkFunction` that 2.0 dropped from the public API. `KafkaSink` with `DeliveryGuarantee.EXACTLY_ONCE` and `FileSink` implement it.

code

java · 10 lines
java
KafkaSink<String> sink = KafkaSink.<String>builder()
    .setBootstrapServers("broker:9092")
    .setRecordSerializer(KafkaRecordSerializationSchema.builder()
        .setTopic("orders-out")
        .setValueSerializationSchema(new SimpleStringSchema())
        .build())
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("orders-out-")
    .setProperty("transaction.timeout.ms", "900000")
    .build();

go deeper

for a junior

Recall the boundary: checkpoints protect Flink's own state, and duplicates in an external system are a separate problem the sink has to solve. Naming a transactional sink as the answer is enough here.

for a middle

Explain the two hooks — stage at the snapshot, commit when the checkpoint completes — and why the commit cannot happen earlier. Be able to name a connector that does it, such as a Kafka sink with an exactly-once delivery guarantee.

for a senior

Walk every crash window, including the pending transaction restored and committed on recovery, and show the operational config that makes it real: transaction timeout above the checkpoint interval, a stable transactional id prefix, committed-read consumers downstream.

for a principal

Own the guarantee as a system property, not a connector setting. Decide where transactional sinks are worth their latency and operational weight versus idempotent writes, and set the standard for how teams state guarantees at pipeline boundaries.

## Where the guarantee stops A completed Flink checkpoint means: every operator's state and every source's position correspond to one consistent cut of the input. If a task fails, Flink restores that state and replays from that cut. Everything inside the job is therefore exactly-once — each input record is *reflected in state* exactly once, no matter how many times it is physically processed. Nothing about this constrains the outside world. If the sink wrote row 4,001 to the database and the job then fails back to a checkpoint taken at record 4,000, the row is still in the database and the replay writes it again. The engine has no way to un-send a network call. So end-to-end exactly-once is a property of the *sink*, and Flink's job is only to give the sink a trustworthy signal about when a cut has become permanent. ## The two-phase-commit contract Flink gives sinks two hooks that map onto a two-phase commit: - **Snapshot** (barrier passes the sink task): the sink flushes everything written so far into a durable but not-yet-visible form, and stores the *handle* to that pending transaction in the checkpoint state. This is the vote-to-commit phase. - **Completion notification** (the coordinator confirms every subtask acknowledged): the sink commits the pending transaction. This is the commit phase, and it is the first instant at which the job knows the snapshot will never be rolled back. In Flink 2.3 these hooks belong to the **Sink V2** API in `org.apache.flink.api.connector.sink2`. A sink that implements `SupportsCommitter` is built from a `CommittingSinkWriter`, whose `prepareCommit()` turns the open transaction into *committables* at the checkpoint, and a `Committer`, whose `commit(...)` makes them visible once the checkpoint completes. The committables are checkpointed on their way to the committer, which is how a restarted job finds them again. The older `TwoPhaseCommitSinkFunction` — `beginTransaction`, `preCommit`, `commit`, `abort`, `recoverAndCommit` — was built on `SinkFunction`, which 2.0 removed from the public API; only a deprecated internal copy survives under `org.apache.flink.streaming.api.functions.sink.legacy`, so new sinks are written against Sink V2. Different code, identical contract. ## Walking the failure cases The contract only earns its keep if it is correct at every crash point: 1. **Crash before pre-commit.** The transaction was never made durable; it is aborted or abandoned. Recovery replays those records and writes them into a fresh transaction. No duplicate, no loss. 2. **Crash between pre-commit and commit.** The checkpoint may or may not have completed. If it completed, the transaction handle was in it: recovery restores the committable and commits it — a second time, if the crash came just after the commit — which is why the Sink V2 `Committer` contract requires commits to be idempotent: a committable that was already committed must leave the external system unchanged and is reported with `signalAlreadyCommitted()`. If the checkpoint did not complete, the job falls back to an earlier one and the pending transaction is aborted. 3. **Crash after commit.** The checkpoint completed and the data is visible; recovery starts from a cut at or after it, so nothing is rewritten. That is the whole argument, and being able to walk it is what a senior interview is testing. ## Kafka as the worked example ```java KafkaSink<String> sink = KafkaSink.<String>builder() .setBootstrapServers("broker:9092") .setRecordSerializer(/* ... */) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix("orders-out-") .build(); ``` Three things fall out of this that candidates routinely miss: - **Consumers must read committed.** Records inside an uncommitted Kafka transaction are on the broker but invisible to a consumer configured for committed reads. A downstream consumer left on the default isolation will see aborted output and the whole exercise is wasted. - **The transaction timeout must exceed the checkpoint interval.** `KafkaSink` opens a fresh transaction right after each checkpoint's snapshot and commits it only when the *next* checkpoint completes, so every transaction lives for roughly a full checkpoint interval, plus any downtime before recovery commits it. Unless you override it, the builder sets the producer's `transaction.timeout.ms` to one hour. If the producer's transaction timeout elapses first, the broker aborts it and that data is *lost*, not duplicated. The producer's timeout must also stay under the broker's maximum, so a long checkpoint interval forces you to raise both. - **The transactional id prefix must be stable and unique per job.** Two jobs sharing a prefix will fence each other's producers. `FileSink` does the same dance without a transaction manager: records go to in-progress part files, a part file that rolls becomes *pending* (bulk formats roll at every checkpoint), and when the checkpoint completes the committer renames pending files to their final visible names. ## The latency you are buying Output becomes visible only when a checkpoint completes, so end-to-end latency is at least the checkpoint interval. A job checkpointing every five minutes with a transactional sink is a five-minute-latency pipeline no matter how fast its operators run. This is the single biggest practical objection to exactly-once sinks and the reason many teams choose an idempotent sink instead. ## The alternative If the sink is naturally idempotent — an upsert keyed by a deterministic business key, an overwrite of a partition, a write whose key encodes the record's identity — then replayed records overwrite rather than duplicate and the *final result* is correct without any transaction. That is at-least-once processing with an idempotent sink, and for many pipelines it is both cheaper and lower-latency than two-phase commit.

  • What goes wrong if the Kafka transaction timeout is shorter than the checkpoint interval?
    The transaction opened at one checkpoint is still waiting for the next one to complete when the broker's timer expires. The broker aborts it, and every record staged inside it is discarded — data loss, not duplication, which is the counter-intuitive part. Set the producer's transaction timeout comfortably above the checkpoint interval plus expected recovery time, and raise the broker's maximum transaction timeout to allow it.
  • Why must a sink's commit step be safe to repeat on a transaction that already committed?
    Because the job can crash after committing but before a newer checkpoint records that fact. On restore it finds the same committable and commits it again. The Sink V2 Committer contract therefore requires idempotent commits: an already-committed committable must leave the external system unchanged and be reported with signalAlreadyCommitted(), otherwise the job fails on every restore and never escapes the loop.
  • A downstream service still sees duplicates from an exactly-once Kafka sink. What is the most likely cause?
    The consumer is not reading with committed isolation, so it sees records from transactions that were later aborted. The sink is doing its job; the reader is looking past the commit boundary. The other common cause is a consumer that commits its own offsets before its side effects are durable, which reintroduces duplicates entirely outside Flink.
  • How does end-to-end latency change when you switch a sink to exactly-once?
    Output stops being visible on write and becomes visible on checkpoint completion, so minimum end-to-end latency rises to roughly the checkpoint interval plus the time to complete a checkpoint. A pipeline checkpointing every two minutes is a two-minute-latency pipeline regardless of operator speed. If that breaks the SLA, the answer is usually an idempotent sink at at-least-once, not a shorter interval.

It is an escrow account: the money is moved and irreversible from the payer's side at signing, but it is only released to the seller once an independent party confirms the deal closed.

saying these in an interview costs you the question

  • Says exactly-once checkpointing alone prevents duplicate output
  • Thinks Flink can retract records already written to a sink
  • Believes committing at snapshot time is safe enough
  • Leaves the Kafka transaction timeout below the checkpoint interval
  • Forgets that consumers must read only committed records

context