Your Kafka Consumer Crashes After Processing the Message but Before Committing the Offset. What Happens Next?
What Kafka at-least-once delivery means when a consumer crashes before committing, and how offsets, idempotency, deduplication, and transactions prevent duplicate effects.
7 min read
The Payment Was Recorded. Then the Consumer Died.
A Kafka consumer receives this event:
{
"eventId": "evt_8f12",
"orderId": "ord_4201",
"type": "PaymentCaptured",
"amount": 12500
}It writes the payment to PostgreSQL and sends a receipt.
Before it commits the Kafka offset, the process crashes.
When the consumer restarts, Kafka delivers the same message again.
Did Kafka fail?
No. This is the expected gap in at-least-once delivery:
process message
|
v
side effect succeeds
|
v
consumer crashes
|
X offset not committed
|
v
message is delivered againWhat the Offset Actually Means
Within a consumer group, Kafka tracks the next position to read for each partition.
Partition 0
offset 40 processed
offset 41 processed
offset 42 delivered now
offset 43 waitingCommitting offset 43 means:
"For this consumer group, resume from offset 43."It does not prove that an email was delivered, a payment row was inserted, or an external API accepted a request. It records Kafka consumption progress.
The Three Timing Windows
Commit After Processing
read -> process -> commit offsetIf the consumer crashes after processing but before the commit, the message is processed again. This gives at-least-once behavior.
Commit Before Processing
read -> commit offset -> processIf the consumer crashes after the commit but before processing, Kafka considers the message consumed. The effect may never happen. This behaves like at-most-once delivery.
Make Effect and Progress Atomic
Kafka and your database are separate systems. A normal database transaction cannot atomically commit a PostgreSQL write and a Kafka consumer offset.
That is why production designs usually make processing idempotent instead of pretending the failure window does not exist.
Idempotency With a Processed-Event Table
Create a unique record for every event:
CREATE TABLE processed_events (
consumer_name text NOT NULL,
event_id text NOT NULL,
processed_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (consumer_name, event_id)
);Insert the deduplication marker and business change in the same transaction:
async function handlePaymentCaptured(event: PaymentCaptured) {
await db.$transaction(async (tx) => {
const inserted = await tx.$executeRaw`
INSERT INTO processed_events (consumer_name, event_id)
VALUES ('billing-payment-captured', ${event.eventId})
ON CONFLICT DO NOTHING
`;
if (inserted === 0) return;
await tx.payment.create({
data: {
orderId: event.orderId,
amount: event.amount,
sourceEventId: event.eventId,
},
});
});
}Then commit the Kafka offset only after the transaction succeeds.
duplicate event
|
v
unique insert conflicts
|
v
handler performs no second business write
|
v
offset can be committed safelyThe deduplication record must share the transaction with the effect. If you write it first in a separate transaction and crash before the business update, the retry will incorrectly skip unfinished work.
Make the Business Operation Naturally Idempotent
Sometimes a separate processed-events table is unnecessary.
If source_event_id is unique on the payment row:
CREATE UNIQUE INDEX payments_source_event_id_key
ON payments (source_event_id);Then duplicate inserts become harmless conflicts.
Likewise, this operation is naturally idempotent:
UPDATE orders
SET payment_status = 'captured'
WHERE id = $1
AND payment_status <> 'captured';But this is not:
UPDATE accounts
SET balance = balance + $1
WHERE id = $2;Replaying it adds the amount twice unless the event identity participates in the write.
External Side Effects Are Harder
Sending an email or calling a payment provider cannot usually share your PostgreSQL transaction.
Use an outbox:
Kafka handler transaction
|
+--> update order
+--> mark event processed
+--> insert receipt_email into outbox
|
COMMIT
|
v
outbox worker sends email with idempotency keyThe outbox makes the intention durable. The worker still needs idempotency because it can also crash after sending but before marking the outbox row complete.
The same reasoning applies to background jobs that fail halfway through.
What Kafka Transactions Can and Cannot Do
Kafka transactions can atomically consume from Kafka and produce to Kafka while committing consumed offsets as part of the transaction.
Kafka input topic
|
v
consumer transforms event
|
v
Kafka output topic + offsets committed atomicallyThis is useful for Kafka-to-Kafka pipelines.
It does not automatically make a write to PostgreSQL, an email provider, or a third-party payment API exactly once. Those systems still require their own idempotency or coordination strategy.
Manual Commit Flow
A simplified consumer loop is:
await consumer.run({
autoCommit: false,
eachMessage: async ({ topic, partition, message }) => {
const event = parseEvent(message.value);
await handleIdempotently(event);
await consumer.commitOffsets([
{
topic,
partition,
offset: (BigInt(message.offset) + 1n).toString(),
},
]);
},
});Real clients also require careful handling of batches, heartbeats, rebalances, and shutdown. Do not copy a commit strategy without understanding the client library’s processing model.
Diagnose Duplicate Processing
Look for:
- Consumer restarts or container evictions
- Group rebalances
- Processing time longer than poll or session settings allow
- Offset commit failures
- Messages repeatedly reaching a dead-letter topic
- The same event ID appearing in logs more than once
- Side effects without stable idempotency keys
Useful metrics include:
consumer_lag
records_processed_total
duplicate_events_total
processing_duration_ms
offset_commit_failures_total
consumer_rebalances_total
dead_letter_totalLog topic, partition, offset, event ID, handler name, and processing result. Without those fields, duplicate processing can look like two unrelated requests.
Poison Messages and Retries
A message that always fails should not block a partition forever.
main topic
|
v
bounded retry topics with delay
|
+--> success
|
v
dead-letter topicRetries need:
- A maximum attempt count
- Backoff
- Failure classification
- Monitoring and replay tooling
- Idempotent handlers
Retrying a permanent schema or validation error indefinitely only creates lag.
Trade-offs
| Approach | Benefit | Cost |
|---|---|---|
| Unique business key | Simple and strong | Not every effect has a natural key |
| Processed-event table | Generic deduplication | Storage and cleanup |
| Outbox | Durable external intent | Worker and operational complexity |
| Kafka transaction | Atomic Kafka read/write | Does not cover arbitrary external systems |
| At-most-once commit | Avoids duplicates | Can lose effects |
Production Best Practices
- Put a globally stable event ID in the event envelope.
- Assume a handler can receive the same event more than once.
- Store the deduplication marker with the business change atomically.
- Commit offsets only after durable processing succeeds.
- Use idempotency keys with external APIs when supported.
- Put external work behind a transactional outbox where appropriate.
- Bound retries and route poison messages visibly.
- Monitor lag, rebalances, duplicates, and commit failures.
- Test crashes at every boundary, not only happy paths.
- Document the exact guarantee of each side effect.
Conclusion
When the consumer crashes after processing but before committing, Kafka redelivers the message because the committed offset still points before it.
The reliable response is not to hope the timing window is rare. It is to make the handler safe under replay:
stable event identity
+
atomic business write
+
idempotent external effects
+
offset committed after successAt-least-once delivery is manageable when duplicates are part of the design rather than a surprise during an incident.
References
Related
Written by
Faisal
Software engineer writing about backend systems, Node.js, system design, scalable applications, and modern web and mobile development.