How does exactly-once delivery for SOURCE connectors work in Kafka Connect, and what is required to enable it?
answer
- KIP-618, Kafka 3.3
- data + offsets in ONE producer transaction
- exactly.once.source.support=enabled on ALL workers
- exactlyOnceSupport() SUPPORTED / exactly.once.support requested
- transaction.boundary = poll|interval|connector + zombie fencing
basics
~20 sConnect (KIP-618, Kafka 3.3+) can run source connectors exactly-once by writing the produced records and their source offsets together in a single Kafka transaction. You enable it cluster-wide with exactly.once.source.support=enabled, and the connector must declare exactly.once.support and define transaction boundaries.
solid answer
~40 sExactly-once source (KIP-618, GA in Kafka 3.3) makes source delivery exactly-once by writing data records AND their source offsets atomically inside one Kafka producer transaction — so offsets can never diverge from the data, eliminating the duplicate window of at-least-once. Requirements: set exactly.once.source.support=enabled on every distributed worker (a one-time, all-workers rolling upgrade), which also gives each task a transactional producer and a per-connector offsets topic. The connector must advertise SourceConnector.exactlyOnceSupport() returning SUPPORTED (or the operator overrides via exactly.once.support=requested in the connector config). Transaction boundaries can be framework-managed (transaction.boundary=poll or interval) or connector-defined (transaction.boundary=connector) via a TransactionContext. Connect also fences out zombie tasks so a stale task can't commit a stale transaction.
go deeper
Know it exists (Kafka 3.3+) and means no duplicates from source connectors when enabled.
Explain data+offsets in one transaction and the worker-level enable flag.
Detail requirements: cluster-wide enable, exactlyOnceSupport(), transaction.boundary modes, per-connector offsets topic.
Reason about zombie fencing, ACL requirements, the rolling-upgrade phases, latency trade-offs, and the source-vs-sink scope boundary.
## The problem it solves By default a source connector is **at-least-once**: it produces records, then later commits source offsets, so a crash between the two replays records = duplicates. **KIP-618** (GA in **Apache Kafka 3.3**) closes that gap by making the produce + offset-commit **atomic**. ## Core mechanism: one transaction for data + offsets Each source task is given a **transactional Kafka producer**. Within a single transaction it writes: 1. the data records to their destination topics, AND 2. the corresponding **source offsets** to the connector's offsets topic. The transaction commits both or neither. Because offsets and data are in the same transaction, they can never diverge — there is no window where data is produced but offsets aren't (or vice versa). Downstream **read_committed** consumers only see records from committed transactions, so duplicates from replays are never visible. ## What you must enable / what's required 1. **Worker config:** `exactly.once.source.support=enabled` on **every** worker in the distributed cluster. This is a two-phase rolling upgrade (you first roll out the code/`preparing` then `enabled`) and applies cluster-wide. Standalone mode does **not** support exactly-once source. 2. **Per-connector offsets topic:** with EOS enabled, each connector gets its own offsets topic (default name derived from the connector, overridable via `offsets.storage.topic` in the connector config) so transactions are isolated per connector. 3. **Connector capability:** the `SourceConnector` must implement `exactlyOnceSupport()` returning `SUPPORTED`. If it returns `UNSUPPORTED` you can still *request* it with connector config `exactly.once.support=requested` (best-effort) vs `required` (fail if unsupported). 4. **Transaction boundaries** via connector config `transaction.boundary`: - `poll` (default) — one transaction per `poll()` batch. - `interval` — commit every `transaction.boundary.interval.ms`. - `connector` — the connector controls boundaries itself using the `TransactionContext` (`transactionContext.commitTransaction()` / `abortTransaction()`), useful when logical units span multiple polls. ## Zombie fencing During rebalances a previous task instance could linger ('zombie') and try to commit a stale transaction. Connect performs **zombie fencing**: the leader bumps producer epochs / fences old transactional IDs so only the current task generation can commit. This requires the worker principal to have the necessary ACLs (e.g. IdempotentWrite, transactional-id permissions). ## Important boundaries - EOS source covers source -> Kafka. It does **not** by itself make a *sink* exactly-once into an external system (that depends on the sink's idempotency / its own preCommit logic). - A connector that can't bound transactions sensibly (e.g. infinite single poll) may not be a good EOS fit. - Enabling EOS slightly raises latency/overhead because of transaction commits and read_committed. ## Quick comparison - At-least-once: produce, then commit offsets separately, on `offset.flush.interval.ms` -> duplicate window. - Exactly-once: produce + commit offsets in one transaction, on `transaction.boundary` -> no duplicate window for read_committed consumers.
- Can you enable exactly-once source on just one worker in a distributed cluster?No. It must be enabled on every worker via a rolling upgrade; mixed settings are unsupported. Standalone mode doesn't support it at all.
- What stops a 'zombie' task from committing a stale transaction after a rebalance?Zombie fencing: Connect fences old transactional IDs / bumps producer epochs so only the current task generation can commit, which requires the worker to hold the right transactional-id ACLs.
- Does enabling EOS source make a JDBC sink connector exactly-once into the database?No. EOS source only guarantees source->Kafka. Exactly-once into an external sink depends on that sink's own idempotency/preCommit handling.
saying these in an interview costs you the question
- Claiming exactly-once source works in standalone mode.
- Saying you can enable it per-connector without the worker-level exactly.once.source.support=enabled.
- Thinking EOS source also guarantees exactly-once into sink targets.
- Forgetting that offsets and data are committed in the SAME transaction (the whole point).