Duplicate records in a downstream sink occur because Kafka operates under at-least-once delivery by default, causing duplicate messages during producer retries, consumer crashes, or group rebalances. If a consumer successfully inserts records into a database table but crashes before committing its partition offset to __consumer_offsets, its replacement consumer will read and insert those same records again. Resolving duplicates requires enabling producer idempotence, committing offsets strictly after sink writes finish, and designing downstream database sinks to perform idempotent upserts keyed by business event IDs.
Common causes of duplicate delivery
Duplicate events usually enter the pipeline at one of three failure points:
- Producer retries without idempotence: A producer writes a record, but a network blip drops the broker acknowledgment. The producer retries and writes the message a second time.
- Consumer crash before offset commit: The consumer polls 100 records, inserts them into a Postgres table, and crashes before calling commitSync(). The newly assigned consumer restarts from the previous committed offset, reinserting the entire batch.
- Rebalance mid-batch: If batch processing exceeds max.poll.interval.ms, the coordinator evicts the consumer and reassigns its partitions while the original consumer is still processing the batch.
Hardening producers and consumers
At the producer boundary, set enable.idempotence=true. The broker will track producer IDs and sequence numbers to discard retried duplicate batches automatically.
At the consumer boundary:
- Disable enable.auto.commit to avoid premature background commits.
- Commit offsets only after the downstream sink write completes successfully.
- For pipelines that stay entirely inside Kafka, use transactional producers and set consumer isolation.level=read_committed for end-to-end exactly-once semantics.
Building idempotent sinks
Because network partitions make true distributed exactly-once writes across external databases difficult, the most reliable defense is making the sink idempotent. Instead of appending raw records with basic insert statements, use primary key idempotent upserts:
INSERT INTO orders (order_id, user_id, amount, updated_at)
VALUES ('ord_101', 'usr_50', 89.99, NOW())
ON CONFLICT (order_id)
DO UPDATE SET amount = EXCLUDED.amount, updated_at = EXCLUDED.updated_at;A short explanation on storage idempotency guarantees:
- If the consumer replays offset 50 due to a crash, the database updates the existing row instead of adding a duplicate row.
- In analytical lakehouses like Iceberg or Delta Lake, use MERGE INTO operations keyed by unique event IDs to deduplicate reprocessed streaming batches.