When would you run a Flink pipeline at-least-once with idempotent sinks instead of exactly-once?
answer
- correctness has two routes, not one
- ask what a duplicate actually does
- visibility waits for the checkpoint
- transactions bring their own failures
- publish the guarantee at the boundary
basics
~20 sWhen the sink is naturally idempotent — keyed upserts, partition overwrites — replayed records overwrite rather than duplicate, so at-least-once gives the same final result without alignment stalls, transaction machinery, or output that is invisible until the next checkpoint completes.
solid answer
~40 sExactly-once is not free: alignment costs latency under load, transactional sinks make output visible only when a checkpoint completes, and transactions bring their own failure modes — timeouts, fenced producers, hanging transactions after an unclean shutdown. Choose it when duplicates are *observable and harmful* and the sink genuinely supports transactions. Choose at-least-once with an idempotent sink when the write is an upsert by a deterministic key, an overwrite of a partition, or a set membership — cases where replaying a record converges on the same value. That combination gives the same end state at lower latency and much lower operational surface. The decision is per pipeline, and it should be recorded as a stated guarantee at the boundary, not left implicit in a connector setting.
code
java · 18 lines// at-least-once route: replay overwrites, no transaction needed
// JdbcSink is org.apache.flink.connector.jdbc.core.datastream.sink.JdbcSink
JdbcSink<Balance> sink = JdbcSink.<Balance>builder()
.withQueryStatement(
"INSERT INTO account_balance (account_id, balance, as_of) VALUES (?, ?, ?) " +
"ON CONFLICT (account_id) DO UPDATE SET balance = EXCLUDED.balance, as_of = EXCLUDED.as_of",
(stmt, b) -> {
stmt.setString(1, b.accountId());
stmt.setBigDecimal(2, b.balance());
stmt.setTimestamp(3, Timestamp.from(b.asOf()));
})
.withExecutionOptions(JdbcExecutionOptions.builder().withBatchSize(500).build())
.buildAtLeastOnce(new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:postgresql://db:5432/ledger")
.withDriverName("org.postgresql.Driver")
.build());
balances.sinkTo(sink);go deeper
Recall that duplicates are only a problem for some sinks: overwriting the same key by a stable identifier is safe under replay, adding to a running total is not.
Explain the two routes to a correct result — transactional commit versus idempotent write — and name the cost of each, particularly that transactional output is not visible until the checkpoint completes.
Argue the choice for a concrete pipeline: sink capability, latency budget, whether duplicates accumulate, and the operational load of transaction timeouts and leftover open transactions.
Set the policy across a fleet: a default guarantee, the criteria for escalating, where deduplication belongs when the sink supports neither route, and the requirement that every pipeline publishes its guarantee and the reader-side conditions at the boundary.
## Framing the decision The question is not "do we want correct results" — both options can be correct. It is which of two routes to correctness fits the sink, the latency budget and the operational appetite. Four dimensions decide it. ## 1. What the sink can actually do Sinks fall into three groups. - **Naturally idempotent.** An upsert keyed by a deterministic business key, an overwrite of a partition or a file with a deterministic name, setting a value rather than incrementing it. Replay overwrites; the final state is the same. These want at-least-once, and paying for transactions buys nothing. - **Transaction-capable.** Kafka with transactions, a relational database, a file sink that renames on commit. These *can* do exactly-once through the two-phase-commit contract, at a price. - **Neither.** A REST call, an email, a payment. No transaction and no idempotence — until you add an idempotency key so the receiver deduplicates. Here the real work is at the receiver, not in Flink. In Flink 2.3 the choice is made per sink, not per job: the JDBC connector's builder ends in `buildAtLeastOnce(...)` or `buildExactlyOnce(...)`, and `KafkaSink` takes a `DeliveryGuarantee`. If the sink is in the first group the argument is essentially over. Most of the interesting cases sit in the second, where both routes are open. ## 2. Latency A two-phase-commit sink makes output visible only at checkpoint completion. Minimum end-to-end latency therefore becomes roughly the checkpoint interval. A pipeline that must surface events within seconds and checkpoints every two minutes cannot use a transactional sink without either shrinking the interval — which raises snapshot cost and, with large state, may not be feasible — or abandoning the latency target. At-least-once with an idempotent sink publishes immediately. For alerting, serving-layer updates and dashboards, that difference usually decides it. ## 3. Whether duplicates are observable Ask what a duplicate actually does downstream: - **Overwritten** — an upsert or a last-value store. Harmless. - **Accumulated** — a running total, an append-only log, a counter. Corrupting, and this is where exactly-once earns its cost. - **Externally visible** — a notification sent twice, a payment taken twice. Unacceptable, and usually solved with an idempotency key at the receiver rather than inside the engine. A useful sharpening question: could a downstream consumer deduplicate more cheaply than the pipeline can guarantee uniqueness? Often it can, and pushing the responsibility there keeps the streaming job simple. ## 4. Operational surface Exactly-once through transactions adds a standing set of failure modes that someone has to own: - transaction timeouts that must be kept above the checkpoint interval, or staged data is aborted and *lost*; - transactional identifiers that must be stable per job and unique across jobs, and that interact with rescaling; - transactions left open by an unclean shutdown, blocking committed readers until they time out; - a hard coupling between output visibility and checkpoint health, so a job whose checkpoints are failing silently stops producing visible output while still looking alive. None is a reason to avoid exactly-once, but all of them are reasons to want it only where it is needed. A fleet where every pipeline runs transactional sinks by default carries that surface everywhere, including on the jobs whose sinks were idempotent all along. ## Deciding, and writing it down A workable rule for a platform: 1. **Default to at-least-once with an idempotent sink.** Design keys so replay converges — this is a modelling decision made at design time, not a runtime setting. 2. **Escalate to exactly-once** when the write is accumulative, the sink supports transactions, and the latency budget can absorb the checkpoint interval. 3. **For non-transactional, non-idempotent side effects**, push deduplication to the receiver with an idempotency key. Never claim exactly-once for a pipeline that ends in an ordinary HTTP call. 4. **Publish the guarantee** at every pipeline boundary — what is guaranteed, and under which failure. The commonest production incident in this space is not a broken guarantee but a *misunderstood* one, where a downstream team assumed uniqueness that nobody ever promised. ## The claim to be sceptical of "We enabled exactly-once, so there are no duplicates" is almost always incomplete. The mode governs Flink's internal state. Output uniqueness requires the sink contract, a reader that respects commit boundaries, and no non-transactional side effects anywhere in the job. A principal-level answer names the whole chain and says explicitly where the guarantee stops.
- What is the cheapest way to turn an accumulating sink into an idempotent one?Stop writing deltas and write the current value instead. A counter that receives "increment by one" is corrupted by replay; the same aggregation written as "the total for this key at this window is N", keyed by the key and window, converges no matter how many times it arrives. The change is in how the record is modelled, not in the engine, and it usually costs far less than adopting transactions.
- How would you state a pipeline's guarantee to a downstream team?Name three things: the delivery property at the boundary (at-least-once, or exactly-once for state and output), the failure under which it holds, and what the consumer must do to see it — for example reading only committed records, or tolerating replayed keys because writes are upserts. A guarantee that does not say what the reader must do is not a guarantee, and the mismatch is the most common source of duplicate-data incidents.
- Is exactly-once ever the wrong choice even when the sink supports transactions?Yes, when latency dominates. Transactional output is invisible until the checkpoint completes, so a sub-second alerting path cannot use it unless the interval shrinks to a point that may be infeasible with large state. It is also wrong when checkpoint health is fragile, because output visibility becomes hostage to snapshot success and a job with failing checkpoints goes quiet while still appearing healthy.
saying these in an interview costs you the question
- Treats exactly-once as strictly better with no cost
- Ignores that transactional output waits for the checkpoint
- Claims exactly-once for a pipeline ending in an HTTP call
- Never asks whether duplicates are actually observable
- Sets the guarantee per cluster instead of per pipeline