A customer gets two confirmation emails. A payment is captured twice. An inventory count drops by two for one sale. The logs show the same message handled twice, a few seconds or a few minutes apart. The queue is not broken. It is doing exactly what it promised: deliver every message at least once.
An SQS or Kafka message processed twice is expected behavior, not a bug in the broker. SQS standard queues and Kafka consumer groups deliver at least once, and any gap between doing the work and acknowledging it can produce a second delivery. Design the consumer so the second delivery is harmless:
- Dedupe on a stable event id, recorded under a unique constraint in the same transaction as the state change.
- Guard state transitions and pass idempotency keys to external APIs.
- Acknowledge after commit: delete the SQS message or commit the Kafka offset last.
- Tune timeouts (visibility timeout,
max.poll.interval.ms) to make duplicates rarer, never to make them impossible.
This guide shows where duplicates come from in SQS and in Kafka, why “exactly-once” features do not cover your database or your payment provider, and how to build a consumer that stays correct when a message is processed twice. If the duplicates you see come from an HTTP provider rather than a queue, How to Prevent Duplicate Webhook Processing covers that side.
Why is a message processed twice?
A message is processed twice because every queue consumer has an ack gap: the time between doing the work and telling the broker the work is done. At-least-once delivery means that if the broker has not received the acknowledgement, it will deliver the message again. Anything that happens inside the gap, such as a crash, a timeout, a deploy, or a rebalance, turns into a redelivery.
You can move the acknowledgement before the work, which is at-most-once: no duplicates, but a crash loses the message. Most systems cannot afford lost orders or lost payments, so they choose at-least-once and accept duplicates. That choice only works if the consumer is idempotent, meaning that processing the same message a second time leaves the system in the same state as processing it once.
The rest of the problem is identifying which event opened the gap in your system, because that tells you how often duplicates will happen. The fix, idempotency, is the same for all of them.
Where do SQS duplicates come from?
In Amazon SQS, a received message is not removed from the queue. It is hidden for the visibility timeout (30 seconds by default, up to 12 hours) and deleted only when the consumer calls DeleteMessage with the receipt handle. If the timeout expires first, the message becomes visible again and any consumer can receive it. The AWS documentation on the SQS visibility timeout describes this lifecycle.
| Cause | How to recognize it | What reduces it |
|---|---|---|
| Processing longer than the visibility timeout | ApproximateReceiveCount of 2 or more on messages that succeeded; duplicates about one timeout apart | Timeout above worst-case processing, smaller batches, and a heartbeat that calls ChangeMessageVisibility |
| Crash or deploy after work, before delete | Duplicates cluster around restarts, OOM kills, or rollouts | Graceful shutdown that stops receiving and finishes in-flight messages |
| Lambda batch with one failure | Whole batches reprocessed when a single record throws | Enable ReportBatchItemFailures and return only the failed ids |
| Standard queue duplicate copy | Rare duplicates with no timeout or crash to explain them | Nothing on the queue side; this is the at-least-once contract |
| Producer resend | Two messages with different MessageIds and the same business payload | A producer-assigned event id; FIFO deduplication for resends within 5 minutes |
AWS Lambda deserves a note of its own. With an SQS event source, Lambda receives a batch, and by default a function error makes the whole batch visible again after the timeout, including the records that succeeded. AWS’s guidance on handling errors for an SQS event source describes partial batch responses, and AWS recommends setting the queue’s visibility timeout to at least six times the function timeout.
def handler(event, context):
failures = []
for record in event["Records"]:
try:
handle(json.loads(record["body"])) # must still be idempotent
except Exception:
log.exception("failed", extra={"message_id": record["messageId"]})
failures.append({"itemIdentifier": record["messageId"]})
return {"batchItemFailures": failures}
FIFO queues change the send side, not the receive side. A FIFO queue drops a second send with the same MessageDeduplicationId (or the same body, with content-based deduplication) within a 5-minute interval. It does not stop a consumer from processing a message twice when the visibility timeout expires or the consumer crashes before deleting it.
Where do Kafka duplicates come from?
A Kafka consumer does not delete messages; it records progress as a committed offset per partition. When a partition changes owner, because a consumer crashed, restarted, or was removed from the group, the new owner starts from the last committed offset. Every record processed after that commit is processed again.
- Crash between processing and commit. With manual commits after the batch, everything in the batch since the last commit is replayed.
- Auto commit.
enable.auto.commitdefaults totrueand commits roughly everyauto.commit.interval.ms(5,000 ms), from insidepoll(). A crash replays up to one interval of processed records. - Rebalance after a slow batch. If the time between two
poll()calls exceedsmax.poll.interval.ms(300,000 ms by default), the consumer is removed from the group and its partitions are reassigned. The new owner replays from the last commit while the old consumer may still be processing. - Upstream resends. The idempotent producer (enabled by default since Kafka 3.0, unless conflicting settings disable it) stops broker-level duplicates from producer retries. It does not stop an application from publishing the same event twice, for example an outbox relay that crashes after sending but before marking the row as sent.
max.poll.interval.ms, and the partition moves before the commit.This is the pattern behind many “Kafka consumer duplicate messages” reports: a downstream slowdown makes batches slower, slow batches trigger rebalances, and each rebalance replays work, which adds load to the same slow dependency. It is a close cousin of a retry storm. The relevant defaults are in the Kafka consumer configuration reference, and the questions developers ask about duplicate messages after a rebalance repeat this sequence over and over.
Doesn’t exactly-once delivery solve this?
Not for the side effects you care about. Both SQS and Kafka offer features described as exactly-once, and both are narrower than the name suggests. Each guarantee covers a specific boundary, and your database writes, emails, and payment calls sit outside it.
| Feature | What it guarantees | What it does not cover |
|---|---|---|
| SQS FIFO deduplication | A second send with the same deduplication id within 5 minutes is accepted but not delivered | Redelivery after a visibility timeout or crash; resends after the 5-minute interval |
| Kafka idempotent producer | Producer retries do not write duplicate records to a partition | The application publishing the same event twice |
Kafka transactions / Streams exactly_once_v2 | Read-process-write to Kafka topics commits output and offsets atomically, for readers using read_committed | Writes to a database, HTTP calls, emails, or any system outside Kafka |
Exactly-once is a property of a boundary. Your payment provider is always on the other side of it.
How do you design a safe consumer?
A safe consumer is one where processing a message twice produces the same outcome as processing it once. You get there by making the consumer’s own writes atomic with a dedupe record, and by making every external effect idempotent on its own. Here is a consumer that is not safe:
while True:
resp = sqs.receive_message(QueueUrl=QUEUE, MaxNumberOfMessages=10, WaitTimeSeconds=20)
for msg in resp.get("Messages", []): # 10 slow messages vs a 30 s lease
order = json.loads(msg["Body"])
payments.charge(order["customer"], order["amount"]) # no idempotency key
mailer.send_receipt(order["email"], order["id"])
db.execute("UPDATE orders SET status = 'paid' WHERE id = %s", (order["id"],))
sqs.delete_message(QueueUrl=QUEUE, ReceiptHandle=msg["ReceiptHandle"])
A crash after charge and before delete_message charges the customer again on redelivery. So does a batch that outlives the visibility timeout. Checking “is the order already paid?” before charging does not fix it either: that is a check-then-act race, and two concurrent deliveries both see pending. (Why “Check Then Insert” Creates Duplicate Records explains that race in detail.)
The safe version follows five steps:
- Give every event a stable id. The producer assigns an event id once, when the business event happens, and carries it in the body or headers.
- Record the id under a unique constraint, in the same database transaction as the consumer’s state change.
- Guard every state transition so it only applies from the expected previous state.
- Make external calls idempotent with a key derived from the business event, or hand them to an outbox worker that does.
- Acknowledge only after commit. Delete the message or commit the offset last.
def handle(event):
with db.transaction() as tx:
first = tx.execute(
"""INSERT INTO processed_messages (consumer, event_id)
VALUES ('billing', %s)
ON CONFLICT DO NOTHING RETURNING event_id""",
(event["event_id"],),
).fetchone()
if first is None:
return # already applied: just acknowledge
tx.execute( # guarded transition
"UPDATE orders SET status = 'charging' WHERE id = %s AND status = 'pending'",
(event["order_id"],),
)
tx.execute( # external effect, handled by a worker
"INSERT INTO outbox (kind, idempotency_key, payload) VALUES ('charge', %s, %s)",
(f"charge:{event['order_id']}", json.dumps(event)),
)
# caller deletes the SQS message only after this returns
The outbox worker calls the payment provider with the stored idempotency key, so a retried charge returns the original result instead of charging again. Providers keep idempotency keys only for a limited window, so the local dedupe record is still the primary guard. Keep the transaction short; making an HTTP call while it is open holds a database connection for the length of the call, which is how a slow provider becomes connection pool exhaustion.
For Kafka, the same handler works with manual commits:
props.put("enable.auto.commit", "false");
props.put("max.poll.records", "100"); // keep a batch well inside max.poll.interval.ms
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
handler.handle(parse(r.value())); // dedupes on the event id, not the offset
}
consumer.commitSync(); // at-least-once: commit only after the work
}
What should the deduplication key be?
The deduplication key should be an event id that the producer assigns once, at the moment the business event happens, and that every copy of the event carries. Broker-assigned identifiers identify a delivery or a record, not the event, so they miss duplicates created before the message reached the broker.
| Key | Catches | Misses |
|---|---|---|
SQS MessageId | Redelivery of the same sent message | Producer resends, which get a new MessageId |
| Kafka topic + partition + offset | Replays after a crash or rebalance | The same event published twice, which gets a new offset |
| Hash of the payload | Byte-identical copies | Copies with a new timestamp or trace id; and it wrongly merges two genuine identical events |
| Producer-assigned event id | Redeliveries, replays, and resends | Only a producer that generates a new id on retry; generate it once and store it with the event |
Two practical notes. Scope the dedupe table by consumer, as in (consumer, event_id), because two consumers of the same topic must each process the event once. And keep dedupe rows longer than any message could be redelivered or replayed: for SQS, at least the retention period (4 days by default, up to 14) plus time spent in a dead-letter queue before redrive; for Kafka, the topic retention or the furthest you would ever reset offsets. If you replay from a dead-letter queue, Dead-Letter Queue Keeps Growing: Where to Start covers doing that safely.
How do you make duplicates rarer?
Tuning cannot make a consumer safe, but it keeps duplicates from becoming a steady load. The goal is that each lease comfortably outlasts the work it covers.
- Visibility timeout above worst-case processing for a whole batch
- Heartbeat with
ChangeMessageVisibilityfor long jobs - Lambda:
ReportBatchItemFailures, and a queue visibility timeout at least 6× the function timeout - Graceful shutdown drains in-flight messages
- Manual commit after processing
- Smaller
max.poll.recordsso a batch fits insidemax.poll.interval.ms - Commit in
onPartitionsRevokedof aConsumerRebalanceListener - Timeouts on downstream calls so one slow call cannot stall a batch
- Log the event id and receive count or offset on every attempt
- Alert on the duplicate rate, not only on errors
- Bound retries and send poison messages to a DLQ
A very long visibility timeout or max.poll.interval.ms also delays recovery: when a consumer really dies, its messages stay hidden, or its partitions stay assigned, until the timeout expires. Size timeouts from measured processing time plus headroom, and use a heartbeat for the long tail rather than one huge value.
How do you test a consumer for duplicates?
Test duplicate delivery directly, because it will not happen by accident in a test environment. Three cases cover most of the risk:
- Sequential duplicate. Deliver the same message twice and assert the state, the outbox, and the external calls are unchanged by the second delivery.
- Crash in the ack gap. Let the handler commit, skip the acknowledgement, and deliver again. The second run must be a no-op.
- Concurrent duplicate. Run two handlers for the same message at the same time, released together by a barrier. Exactly one should apply the change. This is the test that catches a dedupe check written as a
SELECT.
These belong in review as well as in tests. When a pull request adds a consumer or changes one, ask where the ack happens, what runs before it, and what the handler does on a second delivery. The broader checklist is in How to Review a Pull Request for Production Reliability Risks.
How Tomosu helps
Tomosu analyzes a repository and its pull requests for production reliability risks. For queue consumers, the useful evidence is spread across files: the handler, the code that acknowledges, the side effects it calls, and the schema that does or does not back its dedupe logic. For this problem, Tomosu:
- Maps consumer entry points for SQS, Kafka, and similar clients, and where each one deletes, commits, or auto-commits relative to its side effects.
- Flags non-idempotent effects in handlers, such as payment, email, or increment operations with no dedupe record or idempotency key.
- Points at lease risks: slow calls without timeouts inside a batch, and HTTP calls made while a database transaction is open.
- Weighs blast radius, so a consumer that moves money ranks above one that refreshes a cache.
These signals roll up into the Production Reliability Index. To run it on your own consumers, see the repository scan guide.
Scan your repository with Tomosu →
Key takeaways
- An SQS or Kafka message processed twice is normal at-least-once behavior; the consumer has to make it harmless.
- Duplicates come from the ack gap: crashes, expired visibility timeouts, and rebalances between doing work and acknowledging it.
- SQS FIFO deduplication and Kafka exactly-once cover their own boundaries, not your database, emails, or payments.
- Dedupe on a producer-assigned event id, recorded under a unique constraint in the same transaction as the state change.
- Guard state transitions, give external calls idempotency keys, and acknowledge only after commit.
- Tune visibility timeouts, batch sizes, and
max.poll.interval.msto make duplicates rarer, not to make them impossible. - Test sequential, crash-in-the-gap, and concurrent duplicates explicitly.
Frequently asked questions
Why is my SQS message processed twice?
SQS standard queues guarantee at-least-once delivery, so an occasional duplicate copy is expected. The more common cause is the visibility timeout: if a consumer takes longer than the timeout to process and delete a message, the message becomes visible again and another consumer receives it. A crash after processing but before DeleteMessage has the same effect.
Why does a Kafka consumer process the same message twice?
Kafka consumers resume from the last committed offset. If a consumer processes records and then crashes, or loses its partitions in a rebalance before committing, the next owner of the partition reprocesses everything after the last commit. A batch that takes longer than max.poll.interval.ms (5 minutes by default) is a common trigger.
Does an SQS FIFO queue prevent duplicate processing?
It prevents duplicate sends within a 5-minute deduplication interval, using MessageDeduplicationId or content-based deduplication. It does not stop your consumer from processing a message twice if the visibility timeout expires before you delete it, or if you crash after processing. Consumers still need to be idempotent.
Does Kafka exactly-once semantics stop duplicate side effects?
Only inside Kafka. Idempotent producers and transactions give exactly-once results for read-process-write pipelines where the output is another Kafka topic and consumers read with isolation.level=read_committed. Writes to a database, emails, and HTTP calls are outside that transaction, so they still need their own idempotency.
What should I use as the deduplication key?
Use an event id that the producer assigns once, when the business event happens, and carries in the message. The SQS MessageId changes when a producer resends the same event, and a Kafka offset identifies a record rather than an event, so neither catches duplicates created upstream.
How long should I keep processed message ids?
Longer than any message can still be redelivered or replayed. For SQS that means at least the queue’s retention period, which defaults to 4 days and can be up to 14 days, plus time in a dead-letter queue before redrive. For Kafka it means the topic retention or the furthest you would ever reset offsets.
Should I increase the visibility timeout to stop duplicates?
Set it above your realistic worst-case processing time, and extend it with ChangeMessageVisibility for long jobs. That makes duplicates rarer. It does not make them impossible, because crashes, deploys, and upstream resends still cause redelivery, so the consumer must be idempotent either way.
Queues will deliver twice. The question is whether your consumer notices. Tomosu helps you find the handlers that would not, before a duplicate charge does. Assess your repository →