Exactly-Once Semantics in Distributed Systems
“Exactly-once” is the most oversold phrase in distributed systems. Taken literally — a message that crosses a network and is delivered precisely once, no matter what fails — it is impossible, and it has been known to be impossible since the Two Generals problem was written down.
What you can build, and what every system claiming exactly-once is actually building, is at-least-once delivery with idempotent handling. The message may arrive many times; the effect happens once. Since effects are the only thing anyone can observe, that is the property that matters.
Why delivery cannot be exactly-once
A sender transmits and waits for an acknowledgement. The acknowledgement does not arrive. Two things could have happened: the message was lost, or the acknowledgement was lost. From the sender’s position these are indistinguishable — that is the whole problem, and no amount of additional messages resolves it, because each additional message has the same failure mode.
So the sender chooses:
- Retry, and risk delivering twice — at-least-once.
- Give up, and risk delivering zero times — at-most-once.
There is no third option. Every “exactly-once” system picks at-least-once and then does work on the receiving side.
What the receiving side has to do
Three things, and skipping any of them reintroduces duplicates.
Give every message a stable identity. Not a delivery-attempt id, which changes on every retry — a business identity that is the same across all attempts at the same logical operation.
Record that identity atomically with the effect. This is where most implementations fail. If you write the effect and then record the id in a separate step, a crash between the two produces a duplicate on retry. The effect and the record must land in the same transaction.
Bound the record. You cannot remember every id forever. Pick a window comfortably longer than your longest plausible retry, and accept that a duplicate arriving after the window will be processed twice.
The transactional outbox
The hard version of the problem is when the effect and the message live in different systems: write to the database and publish to a queue. There is no transaction spanning both.
The outbox pattern collapses it into one:
BEGIN;
UPDATE accounts SET balance = balance - 100 WHERE id = 42;
INSERT INTO outbox (id, topic, payload)
VALUES ('01HW9Z...', 'payments', '{"account":42,"amount":100}');
COMMIT;
Both rows commit together, so the message exists if and only if the effect happened. A separate relay reads the outbox and publishes. The relay may publish the same row twice — it is at-least-once, like everything else — which is precisely why the consumer still needs its own deduplication.
What Kafka’s transactions actually give you
Kafka’s exactly-once semantics are real, and narrower than the marketing implies. They cover the read-process-write loop within Kafka: consuming from a topic, producing to another, and committing offsets atomically.
They do not extend to side effects outside Kafka. The moment your consumer charges a card or sends an email, you are back to needing an idempotency key and a deduplication store. Kafka made one hop exactly-once; the hop that touches the outside world is still yours to make safe.
Deduplication windows in practice
A dedup store is a cache with correctness consequences, so size it deliberately:
- Too short and a legitimate late retry is processed twice.
- Too long and you are paying to remember ids nobody will ever send again.
Start from your retry policy rather than from a round number. If your producer retries with backoff for up to an hour, a 24-hour window is comfortable and a five-minute one is a bug waiting for a bad afternoon.
When to relax it
Not every operation needs this machinery. The question to ask is what a duplicate actually costs:
- Sending the same push notification twice — mildly annoying.
- Recording the same analytics event twice — a wrong number on a dashboard.
- Charging the same card twice — a support ticket and a refund.
Spend the engineering where the duplicate is expensive, and let the cheap ones be at-least-once with no ceremony. A system that treats every message as if it were a payment is a system nobody can afford to operate.
Quick answers
- Is exactly-once delivery possible?
- No. Across an unreliable network a sender cannot distinguish a lost message from a lost acknowledgement, so it must either risk losing it or risk sending it twice. What is achievable is exactly-once effects, via at-least-once delivery plus idempotent handling.
- What is the difference between exactly-once delivery and exactly-once processing?
- Delivery is about how many times a message arrives, which cannot be bounded to one. Processing is about how many times its effect is applied, which can be bounded to one by deduplicating on a stable id.
- How does Kafka claim exactly-once semantics?
- Through idempotent producers and transactions that atomically commit both the output records and the consumer offsets. It holds within Kafka; the moment a side effect leaves Kafka, such as charging a card, the guarantee no longer covers it.
- How long should a deduplication window be?
- Longer than the maximum plausible retry and redelivery interval, including broker retention and consumer downtime. Too short and a late duplicate slips through; too long and the dedupe store becomes its own scaling problem.
References
Related Discoveries
Lumi's weekly note
A short email when we publish something useful. No spam, unsubscribe anytime.