Blog

Exactly-Once Kafka Sinks: Where Duplicate Writes Begin

Table of Contents

Table of Contents

The sink connector was configured for exactly-once delivery. The target still contained two rows for one payment event.

The incident was anonymous and illustrative, but the sequence is ordinary. A sink task read a record, sent it to the destination, and received a success response. Before it committed the Kafka offset, the worker stopped. After restart, the task read the record again. The destination treated the retry as new because the connector had not given it a durable way to recognize the first attempt.

That is the boundary behind many claims about Kafka sink connectors. Exactly-once behavior inside Kafka can make a record position and a broker-side transaction line up. It does not automatically make an external key-value store, database, or object store apply a side effect once. The useful question is, “Which commit does that claim cover, and what does the target do when the write result is ambiguous?”

1The duplicate key that should not exist

Start with the sequence instead of the marketing label:

  1. The task fetches record orders-3@742 from Kafka.
  2. It maps the event to a destination write.
  3. The destination accepts the write.
  4. The task fails before its source position is committed.
  5. A replacement task reprocesses orders-3@742.

The second delivery is not evidence that Kafka lost the record. It shows that the task did not have one atomic commit spanning Kafka’s progress and the destination’s side effect. Retrying after an uncertain result preserves data, but safe replay requires the destination to recognize the event as the same event.

The failure can happen with no visible error at the target. A timeout may arrive after the remote system has applied the write, or a process may be killed after a database transaction commits but before the connector records progress. A connector that creates a fresh destination ID for every call turns either case into a second object, row, or document.

The target’s write contract therefore matters as much as the connector’s retry policy. A deterministic overwrite may be safe. A conditional create may be safe when the condition is durable. An insert into a table with a unique event key may be safe. A generated identifier, an append-only object name, or a dedupe check performed in a separate transaction leaves an opening for duplicates.

A Kafka sink delivery path showing the broker offset commit and target write as separate boundaries

The phrase “exactly once” needs a subject and a boundary. Kafka’s semantics documentation explains how transactions and offset handling can provide exactly-once processing for a Kafka-to-Kafka flow when the operations share a transaction. An external sink usually cannot join it. Kafka Connect’s connector documentation describes the worker model, but the final guarantee depends on the connector implementation and destination API.

2Where exactly-once leaves the building

A sink has at least two progress records: the source position and the target state. Kafka stores the source position through the Connect worker’s internal state. The target stores a row, item, object, or other effect. These records can be ordered, but they are not automatically one commit.

The safe default ordering is target first, source position second:

Apply the target-side effect, make the result durable, then commit the source position.

If the task crashes between those operations, Kafka replays the record. That is acceptable when the target operation is idempotent. If the task commits the source position first, a crash before the target write can create a gap that replay will not repair. The choice is therefore usually at-least-once delivery plus idempotent application, unless a particular connector and target provide a stronger end-to-end protocol.

There are four different claims teams often compress into one label:

ClaimWhat it coversWhat it does not prove
Exactly-once source processingA source task and Kafka offsets follow a transaction-aware protocolAn arbitrary external target write
Exactly-once Kafka-to-Kafka processingConsumed records, produced records, and offsets share Kafka transaction scopeA database, API, or object store outside Kafka
Idempotent sink writeRepeating the same event produces the same accepted stateCorrect event identity or retention of dedupe state
End-to-end exactly-once effectThe source, connector, target, and recovery protocol share a tested atomic boundaryA stronger guarantee after the system boundary changes

This distinction explains why an offset dashboard can look healthy while a destination is wrong. Offset lag measures source progress; it does not reconcile target state. Review the interval after the target acknowledges and before Kafka records progress, then ask whether a timeout retry preserves the same idempotency key and how an operator recognizes a replay.

3Three sink profiles, three failure surfaces

The destination type changes the shape of the gap. The same Kafka retry can be harmless for one sink and create a second business record for another.

3.1Key-value stores: deterministic keys are the first control

A key-value sink can make a replay safe when the destination key is derived from a stable event identity and the operation replaces the value at that key. A retry of the same event then writes the same logical item. The design still needs a policy for conflicting payloads: if the same key arrives with a different value, silently overwriting it can hide an upstream bug.

Conditional writes add a useful guard. A “put if absent” operation can ensure that the first accepted event owns the key, while a replay returns a condition failure that the connector treats as an already-applied result. The condition must be part of the target’s atomic operation. A read followed by a separate write is a race between workers.

The hard case is an ambiguous timeout. The client must retry with the same key because it cannot tell whether the first write completed. A generated destination ID defeats this design. If dedupe state expires, define the replay window that remains safe. The DynamoDB condition expression documentation is one example of a target-side conditional write contract; other key-value systems use different names and limits.

3.2Transactional databases: put the marker and the row in one transaction

A transactional database gives a sink more room to build a durable boundary. The business row can use an upsert keyed by the event’s identity, or the database can hold a separate dedupe table with a unique constraint. The strongest common pattern is to insert the dedupe marker and apply the business change in the same database transaction.

On a first delivery, the marker insert succeeds and the business mutation commits with it. On a replay, the unique constraint rejects the marker insert, so the sink can treat the event as already applied. If the process fails before the transaction commits, neither the marker nor the business change is visible and a retry can attempt the work again. If it fails after commit but before the Kafka offset commit, the replay finds the marker.

The tables must share a transaction boundary. Writing the business row first and the dedupe marker later reintroduces the gap. So does keeping the marker in an evictable cache or using a non-unique lookup that two tasks can pass at the same time. PostgreSQL’s INSERT ... ON CONFLICT documents the primitive, but the connector must place the marker, business mutation, and error handling in one transaction.

A marker table needs a retention policy, indexes, and a response to late replays. It must record enough payload identity to detect a reused key with changed content. That cost belongs in the delivery design, not inside a claim that the sink “supports transactions.”

3.3Object stores: names and manifests replace row-level dedupe

Object stores do not provide a row-level upsert. A sink that writes every delivery to a generated object name creates a second object when Kafka replays the record. The storage system may be behaving exactly as requested.

The first design option is a deterministic object key. If the event is a complete immutable object, rewriting the same key can be idempotent when the content is deterministic and replacement is acceptable. The second is a conditional create, using the object store’s precondition support to accept the first writer and reject later attempts. The Amazon S3 PutObject documentation describes conditional request behavior; an implementation must verify which preconditions its chosen object store supports.

Batch sinks need another layer. If one object contains many Kafka records, a stable object key does not say which records are present after a retry. Use a manifest or committed batch marker with the source identity range and content checksum, then make readers consume committed manifests rather than every temporary object.

Object storage fits immutable event archives, but reconciliation belongs to the manifest and reader path. An object list is not an idempotency protocol. If the reader cannot tell which duplicate object is authoritative, the sink has moved the ambiguity downstream.

A comparison of key-value, transactional database, and object-store sink failure profiles

  • Key-value: use a stable key and an atomic overwrite or conditional create; inspect conflicting payloads.
  • Transactional database: commit the dedupe marker and business mutation together; make the marker unique and durable.
  • Object store: use deterministic names or committed manifests; do not treat generated filenames as event identity.

The sink’s target system is part of the delivery guarantee. Connector settings cannot erase that fact.

4Idempotency keys and dedupe tables that survive replay

The idempotency key is the join point between a Kafka replay and an already-applied target effect. It must remain the same when the record is retried, moved to another task, or processed after a worker restart.

Start with the strongest identity the event already carries. A producer-assigned event ID travels with the record. If the source has no stable ID, a namespaced source position such as cluster/topic/partition/offset can identify one record within that log. It does not identify a business event republished to a new topic, mirrored into another cluster, or reconstructed with a different offset. Write down that boundary instead of treating offsets as universal identity.

A sound key design answers five questions:

  1. Scope: Is the key unique across tenants, topics, source clusters, and sink instances?
  2. Replay: Does a retry preserve the exact bytes or canonical identity used to build the key?
  3. Conflict: What happens when the same key arrives with a different payload hash or event version?
  4. Retention: How long must the dedupe record remain to cover operational replay and backfill windows?
  5. Migration: Can two connector versions compute the same key during a rolling upgrade?

For a database sink, a minimal dedupe table might look like this:

sql
CREATE TABLE sink_dedupe ( sink_name text NOT NULL, idempotency_key text NOT NULL, payload_hash text NOT NULL, first_seen_at timestamptz NOT NULL, PRIMARY KEY (sink_name, idempotency_key) );

The write path runs inside one transaction. Insert the marker with a unique constraint, inspect whether it created a row, and apply the business mutation only for a new marker. If the marker exists with the same payload hash, acknowledge the replay as already applied. If it exists with a different hash, stop or route the record to a controlled error path. That is a data contract violation, not a normal retry.

The marker can be combined with the business table when a unique business key already expresses event identity. A separate table helps when one event updates several tables, business keys can change, or the sink needs an audit trail for suppressed replays. Either way, the marker and effect must share the target’s durable commit boundary.

Retries, backfills, and dead-letter handling belong in the same design. A dead-letter topic preserves a record that could not be applied; it does not make a non-idempotent target safe. A backfill that reuses the original event ID can be deduplicated, while a re-export with new IDs is a new delivery stream. The Kafka Connect error-handling model describes the retry path; the target contract decides whether a repeated call is safe.

5What changes when the source is AutoMQ?

The source platform can change without moving this boundary. AutoMQ is a Kafka-compatible streaming platform, so a Kafka Connect deployment can use its Kafka-facing interface subject to the connector and workload tests. AutoMQ’s Kafka compatibility documentation describes that source-side boundary.

It does not generate an idempotency key for a destination, include a Kafka offset in a database transaction, or turn an append-only object path into a conditional write. AutoMQ’s architecture overview describes the source data plane; target-side controls still belong to the connector and destination.

That division of responsibility is healthy. A platform should be judged on its source guarantees, while the sink should be judged on its target protocol. Teams comparing Kafka-compatible options can use the published Kafka Connect and AutoMQ data pipeline example for connectivity context. Connectivity is the beginning of the test, not proof of end-to-end delivery semantics. For replay and cutover questions, the blue-green Kafka migration guide keeps acknowledged records and rollback in view.

6A delivery-semantics self-assessment

Run this assessment against one real sink and one representative failure drill. Record evidence beside each answer.

QuestionEvidence to collectIf the answer is unclear
What event identity does a retry carry?Key derivation rule and sample recordsTreat the sink as duplicate-prone
When is source progress committed?Worker and connector offset behaviorAssume a crash can replay the record
Can a target timeout hide a completed write?API contract and injected timeout resultRetry only with the same key
Are the business effect and dedupe marker atomic?Target transaction or conditional-write proofKeep the guarantee at-least-once
What happens on a key collision with changed content?Payload hash, version check, and error pathQuarantine the record for review
How long does dedupe state live?Retention policy tied to replay and backfill windowsDocument the unsafe replay window
Can an operator see suppressed replays?Metrics, logs, audit records, and alertsAdd a replay counter before relying on it
Has the crash window been tested?Failure injection after target success and before offset commitDo not claim exactly-once effects

Eight questions for checking Kafka sink delivery semantics before claiming exactly-once effects

The result should be a precise label. “At-least-once source delivery with idempotent target application” is a useful outcome. “Exactly once” is meaningful only when the source position, target effect, and recovery protocol share a tested boundary. If the target can receive a duplicate but the resulting state is unchanged, call that effectively once and document the conditions that make it true.

The anonymous duplicate row at the start of this article was not fixed by changing a broker acknowledgment setting. The repair was to make event identity durable, put the business write and dedupe decision in the same target-side boundary, and test the restart path that had been missing from the original claim. That is the level at which a sink’s delivery semantics become real.

If you are evaluating a Kafka-compatible source for this workload, start an AutoMQ evaluation with the sink’s failure matrix beside the connectivity test. The source can preserve Kafka’s replayable record path. Your target still needs a write protocol that makes replay safe.

Newsletter

Subscribe for the latest on cloud-native streaming data infrastructure, product launches, technical insights, and efficiency optimizations from the AutoMQ team.

Join developers worldwide who leverage AutoMQ's Apache 2.0 licensed platform to simplify streaming data infra. No spam, just actionable content.

I'm not a robot
reCAPTCHA

Never submit confidential or sensitive data (API keys, passwords, credit card numbers, or personal identification information) through this form.