As a platform architect, how would you reason about end-to-end exactly-once across a pipeline that ingests from an external source, processes in Kafka, and lands in an external sink?
answer
- three boundaries: source / internal / sink
- only Kafka-internal hop is true EOS
- KIP-618 Connect EOS source ≠ external read atomicity
- thread a deterministic id end-to-end, dedupe everywhere
- choose guarantee per data path (cost vs. need)
basics
~20 sTreat the pipeline as three hops: source→Kafka, Kafka→Kafka, Kafka→sink. Only the Kafka-internal hop can be true EOS. The source and sink boundaries are at-least-once unless you add idempotency, dedupe keys, or store-local atomicity. End-to-end exactly-once is engineered at the edges, not granted by Kafka.
solid answer
~50 sI decompose the pipeline into boundaries and ask, at each boundary, whether an atomic commit is possible. The Kafka-internal stage (consume-transform-produce) can be exactly_once_v2 because offsets and outputs share one cluster. The two edges can't: a source connector that reads an external system and produces to Kafka can re-produce on restart (Connect EOS / KIP-618 helps make produce+Connect-offset atomic on the Kafka side, but the *external* read position is the source's problem), and a sink re-applies on redelivery. So I make edges carry deterministic identifiers end-to-end and enforce idempotency: stable record ids from the source, dedupe in stream processing keyed on those ids, and idempotent or offset-tracking writes at the sink (outbox/inbox, upserts, store-local offset). I also pick the guarantee per requirement — financial flows justify the cost of dedupe ledgers; analytics may accept at-least-once. The architecture decision is: exactly-once *effect* via idempotency, not exactly-once *delivery* across systems.
go deeper
Understand that only the Kafka-to-Kafka step can be exactly-once; the edges need idempotency.
Break the pipeline into source/internal/sink and identify which boundary Kafka EOS covers.
Design end-to-end idempotency with a deterministic id and explain Connect EOS limits at the source.
Drive per-data-path guarantee decisions, retention/ordering trade-offs, and an org pattern (deterministic id + idempotent edges) for exactly-once effect.
## Decompose into boundaries End-to-end pipelines cross **three boundaries**, and exactness must be reasoned about at each: 1. **External source → Kafka** (e.g., a Connect source connector, CDC, an HTTP poller). 2. **Kafka → Kafka** (stream processing: consume-transform-produce). 3. **Kafka → external sink** (DB, index, API, object store). Kafka's EOS guarantee covers **only boundary 2**, and only when it's all in one cluster. ## Boundary 1 — source ingestion A source reads an external system and **produces to Kafka**. Two failure gaps: - It may **re-read** the external system after a crash (its committed external read-position vs. produce isn't atomic). - It may **re-produce** to Kafka. **Kafka Connect exactly-once source support (KIP-618)** makes the *produce to Kafka* + *Connect source-offset commit* atomic via Kafka transactions, removing Kafka-side duplicates. But the **external read position** is still the connector's responsibility; whether the source can resume at exactly the right external position (and not re-emit) depends on the source system (CDC log offsets are good; a non-transactional REST poll is not). So ingestion is **at-least-once into Kafka** unless the source supports atomic position+emit. ## Boundary 2 — Kafka-internal processing `exactly_once_v2` (Kafka Streams) or manual transactions bind output records + input offset commit atomically in the cluster. This is the **only true EOS hop**. Requires downstream `isolation.level=read_committed`. ## Boundary 3 — sink Same two-non-atomic-steps issue as any sink: apply effect + commit offset aren't atomic across systems. Achieve **effectively-once** via idempotent upserts, dedupe ledgers (inbox), or storing the Kafka offset in the sink store transactionally. ## The architect's move: end-to-end idempotency Because no protocol spans all three systems, **engineer exactly-once *effect* with deterministic identifiers**: - Assign a **stable, deterministic id** at the source (CDC LSN, business key, or producer UUID) and propagate it unchanged through every hop. - In stream processing, **dedupe** on that id (e.g., a state-store seen-set with retention) to absorb source-side duplicates. - At the sink, make writes **idempotent on that id** (upsert / unique constraint / conditional put) or keep an **inbox** of processed ids. Now a duplicate anywhere collapses to a single effect. ## Choosing the guarantee per requirement Exactness has cost (latency, throughput, storage for dedupe state, operational complexity of transactions/Connect EOS). Decide per data path: - **Financial / billing / ledgers** → invest in end-to-end idempotency + outbox/inbox. - **Analytics / metrics / logs** → at-least-once (or even at-most-once) is often fine; dedupe in the warehouse if needed. ## Failure-mode and operational considerations - **Dedupe state retention**: seen-sets/inbox tables can't grow forever; bound by time/window and accept that very-late duplicates beyond the window may slip through. - **Ordering**: per-partition only; cross-partition or cross-topic global order isn't guaranteed. - **Cross-cluster** (MM2) adds another at-least-once boundary with approximate offset translation. - **Transactional throughput**: EOS transactions add commit overhead; size `transaction.timeout.ms` and batch appropriately. ## One-liner Kafka gives you one exactly-once hop; the architect builds end-to-end exactly-once *effect* by threading a deterministic id through every boundary and enforcing idempotency at the edges — then chooses how much of that to pay for per data path.
- Does Kafka Connect's exactly-once source support (KIP-618) give you exactly-once ingestion from any external source?No. It makes the produce-to-Kafka and the Connect source-offset commit atomic, eliminating Kafka-side duplicates. But resuming at the exact external read position without re-emitting is the source system's responsibility — feasible for log-based CDC, not for a non-transactional REST poll. So ingestion is exactly-once into Kafka only when the source can atomically track its own position.
- How do you keep end-to-end dedupe state from growing unbounded?Bound the dedupe window by time or by a watermark (e.g., retain seen-ids for a max expected duplicate delay), evict older entries, and accept that duplicates arriving beyond that window are rare and handled by other safeguards. Trade window size against memory/storage and the probability of late duplicates.
- When would you deliberately accept at-least-once or at-most-once instead of building exactly-once?For high-volume, low-value-per-record streams like metrics, logs, or clickstream analytics where occasional duplicates or rare losses don't change decisions, and where the cost of dedupe state, transactions, and operational complexity outweighs the benefit. You can also dedupe downstream in the warehouse if needed.
saying these in an interview costs you the question
- Asserting that turning on EOS/exactly_once_v2 makes the whole pipeline exactly-once end-to-end
- Treating KIP-618 Connect source EOS as solving the external read-position problem
- Designing one global guarantee instead of per-data-path guarantees
- Building unbounded dedupe state with no retention strategy