Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. End-to-end exactly-once

Batch & Streaming · Streaming Design

End-to-end exactly-once

Hardbatch-streaming-37
exactly-onceidempotencykafkatransactions

Question

What does end-to-end exactly-once actually require?

Solution

End-to-end exactly-once processing ensures that every event affects downstream outputs and internal system state as if it was processed exactly once, even during task failures and restarts. Achieving this guarantee requires a replayable message source with offset tracking, atomic checkpointing of offsets alongside operator state, and a destination sink that is either transactional or idempotent. Any non-idempotent side effect along the pipeline shatters this guarantee, which is why most production platforms achieve effectively-once outcomes through idempotent writes.

The three architectural prerequisites

To deliver exact processing across failure boundaries, three components must coordinate:

  • Replayable source: Message systems like Apache Kafka or AWS Kinesis allow consumers to rewind and re-read historical offsets upon worker recovery. If a source cannot replay uncommitted data, missing records cannot be reconstructed.
  • Coordinated checkpointing: The stream processor, such as Apache Flink, periodically takes distributed snapshots that record source offsets and internal operator states at consistent logical moments.
  • Coordinated sink: The output sink must participate in the commit protocol. In transactional approaches, sinks implement two-phase commit protocols or Kafka transactions, staging data until the coordinator confirms the checkpoint succeeded.

Transactional sinks versus idempotent writes

Two-phase commit sinks keep transactions open until global checkpoints complete. If an external coordinator times out or fails, long-running transactions can block downstream consumers or leave orphaned data in external systems.

Because two-phase commits introduce operational complexity, experienced data engineers prefer idempotent writes:

At-least-once streaming delivery + Idempotent destination sink
= Effectively-once processing result

Writing records using an UPSERT statement or a Delta Lake MERGE keyed on a unique transaction ID means that reprocessing duplicate records produces identical destination rows without duplicate count increments.

The danger of external side effects

A common junior engineer mistake is triggering external network calls within streaming operators, such as sending emails, publishing push notifications, or calling third-party payment APIs.

Because stream recovery replays events from previous checkpoints, any external HTTP call executed inside a transformation function will fire repeatedly during recovery. External APIs cannot participate in two-phase commit rollbacks, breaking the exactly-once promise for user-facing systems.

PreviousNext