Kafka "at least once": making a consumer idempotent

In the article on event-driven with Kafka, one row of the table sums up a trade-off in a few words: "at least once (the consumer must handle duplicates)". That parenthesis hides half the work of building a system on Kafka. This article explains where duplicates come from, why the "exactly-once" option doesn't remove all of them, and how to write a consumer that absorbs them without side effects.
Where duplicates come from#
A Kafka consumer does two separate things: it processes a message, then it commits that message's offset so it won't read it again. Those two actions aren't atomic. Everything happens in the gap between them.
read the message ──▶ process (write to DB, call an API) ──▶ commit the offset
▲
crash, rebalance or timeout here:
processing is done, the offset isn't committed,
the next consumer will read the message again
If you flip the order (commit first, process after), you get the opposite: a crash in between loses the message. That's "at most once". Between losing a message and processing it twice, most systems pick the duplicate, because a duplicate can be detected and a loss is invisible.
The concrete causes of a re-read are mundane and frequent:
- a redeploy that stops the pod between processing and commit;
- a consumer group rebalance, for example because processing exceeded
max.poll.interval.ms(5 minutes by default); - a producer resending a message after a network timeout although the broker did receive it (on the producer side, this case is handled by producer idempotence, see below).
In other words, duplicates aren't an edge case. On a system that has been running for a few months with regular deploys, they are a certainty.
What "exactly-once" actually guarantees#
Kafka offers two mechanisms often presented as the solution.
The idempotent producer (enable.idempotence=true, on by default since
Kafka 3.0). The broker assigns the producer an ID and each message a
sequence number. If the producer resends a message after a timeout, the
broker recognizes the duplicate and drops it. That fixes duplicates when
writing to Kafka, not when reading.
Transactions (transactional.id on the producer,
isolation.level=read_committed on the consumer). They let you write
messages to several topics and commit the consumed offsets in one atomic
transaction. That's what Kafka Streams uses to provide real exactly-once.
The limit is in the definition: these guarantees cover exchanges between Kafka topics. The "read a topic, transform, write to another topic" pattern can be exactly-once. The "read a topic, write to PostgreSQL, send an email" pattern can't: neither the database nor the mail server takes part in the Kafka transaction.
| Scenario | Exactly-once possible with Kafka alone? |
|---|---|
| Topic → transformation → topic | Yes, with transactions (or Kafka Streams) |
| Topic → database write | No, the database isn't in the transaction |
| Topic → external API call | No |
| Topic → email or SMS | No, and it's the case users notice most |
For every case in the bottom part of the table, which covers most real consumers, idempotence is the consumer's job.
What idempotent means, exactly#
An operation is idempotent when applying it twice leaves the same state as
applying it once. balance = 100 is idempotent. balance = balance + 100
isn't. The whole technique is about turning the second kind of operation
into the first, or detecting that the message was already processed.
There are three ways to do it, from simplest to most general.
1. Make the operation naturally idempotent#
Sometimes changing the shape of the write is enough. An
OrderStatusChanged(orderId=42, status=SHIPPED) event can be applied as many
times as you like: you set the status, you don't increment it. An upsert on
a business key has the same property:
INSERT INTO order_status (order_id, status, updated_at)
VALUES (:orderId, :status, :timestamp)
ON CONFLICT (order_id) DO UPDATE
SET status = EXCLUDED.status, updated_at = EXCLUDED.updated_at
WHERE order_status.updated_at < EXCLUDED.updated_at;The WHERE clause on the timestamp also protects against an old message
re-read after a newer one. When it's possible, this is the best solution: no
technical table, no deduplication logic.
2. A processed-messages table, in the same transaction#
When the operation isn't idempotent by nature (crediting an account, creating an invoice), you record the ID of each processed message, in the same database transaction as the business effect. If the message comes back, the insert hits the unique constraint and you know it was already processed.
CREATE TABLE processed_message (
message_id VARCHAR(64) PRIMARY KEY,
processed_at TIMESTAMP NOT NULL
);@Component
class PaymentReceivedConsumer(
private val jdbc: JdbcTemplate,
private val accounts: AccountRepository,
private val tx: TransactionTemplate,
) {
@KafkaListener(topics = ["payments"], groupId = "accounting")
fun onMessage(event: PaymentReceived) {
tx.executeWithoutResult {
val isNew = jdbc.update(
"""
INSERT INTO processed_message (message_id, processed_at)
VALUES (?, now())
ON CONFLICT (message_id) DO NOTHING
""".trimIndent(),
event.paymentId,
) == 1
if (isNew) {
accounts.credit(event.accountId, event.amount)
}
}
}
}Three details matter here.
The ID comes from the producer, not from Kafka. You could use the
topic, partition and offset triple, but it changes if the message is
republished (replay from another source, cluster migration). A business ID
(paymentId) or a UUID generated when the event is created stays stable.
The insert and the effect share one transaction. If you insert the ID, commit, then credit the account in a second transaction, a crash in between marks the message as processed while the credit never happened. You've swapped a duplicate for a loss.
The table grows. It needs purging, keeping a window longer than the topic's retention: a message can't come back once Kafka has deleted it.
3. External effects: idempotency key or outbox#
For an API call or an email, the local database is no longer enough: you can't unsend an email if the transaction fails afterwards.
Two options, depending on what the called API allows.
The API accepts an idempotency key. Many payment and messaging APIs do
(an Idempotency-Key header). You pass the message ID, and the called
service ignores the second call. It's the simplest solution; just check how
long the service keeps the keys.
The API doesn't accept one. You separate the decision from the execution: the consumer writes to a "to send" table, in the same transaction that marks the message (method 2). A second process reads that table, sends, and marks the row as sent. The Kafka duplicate is absorbed by the table, and the only risk left is a double send on a crash during the send itself, which you can shrink but not fully remove without help from the called service.
The reverse problem: publishing reliably#
Duplicates have a symmetric cause on the producer side. A service that writes to its database and then publishes an event has the same atomicity problem:
// Fragile: if publishing fails, the order exists
// but nobody is notified. If the database fails after
// publishing, we announce an order that doesn't exist.
orders.save(order)
kafkaTemplate.send("orders", OrderCreated(order.id))The transactional outbox pattern fixes this: you write the event to an
outbox table in the same transaction as the order, and a separate process
(a job reading the table, or Debezium reading the database log) then
publishes to Kafka. Publishing can be replayed, so it sometimes produces
duplicates, which brings us back to the idempotent consumer. The two
patterns go together: the outbox guarantees no event is lost, idempotence
guarantees a duplicated event does no harm.
Testing idempotence#
One test is enough to prove the essentials: send the same message twice and check the final state.
@Test
fun `a payment received twice credits the account only once`() {
val event = PaymentReceived(paymentId = "p-123", accountId = "c-1", amount = 50.euros)
consumer.onMessage(event)
consumer.onMessage(event)
assertThat(accounts.balance("c-1")).isEqualTo(50.euros)
}This test is cheap and catches most regressions. The harder case, a crash in the middle of processing, is tested by injecting an exception after the business effect and before the end of the transaction, then replaying the message: the state must be the same as if it had been processed once.
What to take away#
"At least once" isn't a misconfiguration you fix by switching an option on. It's the normal contract of a consumer with effects outside Kafka. Kafka's exactly-once is real, but it stops at the broker's boundary. Beyond it, three tools cover the vast majority of cases: a naturally idempotent write when possible, a processed-messages table in the same transaction otherwise, and an idempotency key or an outbox for external effects. The question to ask for every new consumer fits in one line: what happens if this message arrives twice?


