Company
About Tomosu
Platform
Platform & Agents Indexes How it works Solutions Pricing
Get Started
MCP Server VS Code — Plugin Installation Scan Your Repo — Guide Integrations · GitHub App Integrations · CodeRabbit MCP FAQ
Free Tools
Governance Impact
Resources
Blogs News Download / Free Trial Book a call →
Production Debugging · Queues

SQS or Kafka Message Processed Twice: Designing Safe Consumers

Tomosu AI·14 min read·

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.

Quick answer

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:

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 ACK GAP: WHERE A SECOND DELIVERY COMES FROM Receive poll / ReceiveMessage Process parse, validate, compute Side effects DB write, charge, email Acknowledge delete / commit offset THE ACK GAP: WORK DONE, BROKER NOT YET TOLD Crash before side effects Redelivered, nothing happened yet: harmless. Crash after side effects, before ack Redelivered: every side effect runs again. Processing outlives the lease A second consumer gets it while the first runs. Ack moved before the work No duplicates, but a crash loses the message. The lease is the SQS visibility timeout, or the Kafka group membership bounded by max.poll.interval.ms.
Every at-least-once consumer has a gap between finishing the work and acknowledging it. Duplicates come from that gap.

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.

VISIBILITY TIMEOUT 30 S, PROCESSING 45 S CONSUMER 1 QUEUE CONSUMER 2 processing a batch of 10 (45 s) invisible (visibility timeout) processes the same message DeleteMessage visible again, received 0 s 15 s 30 s 45 s 60 s Result: two consumers ran the side effects, partly at the same time. The later delete does not undo either.
The visibility timeout is a lease, not a lock. When the lease runs out, SQS hands the message to someone else.
CauseHow to recognize itWhat reduces it
Processing longer than the visibility timeoutApproximateReceiveCount of 2 or more on messages that succeeded; duplicates about one timeout apartTimeout above worst-case processing, smaller batches, and a heartbeat that calls ChangeMessageVisibility
Crash or deploy after work, before deleteDuplicates cluster around restarts, OOM kills, or rolloutsGraceful shutdown that stops receiving and finishes in-flight messages
Lambda batch with one failureWhole batches reprocessed when a single record throwsEnable ReportBatchItemFailures and return only the failed ids
Standard queue duplicate copyRare duplicates with no timeout or crash to explain themNothing on the queue side; this is the at-least-once contract
Producer resendTwo messages with different MessageIds and the same business payloadA 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.

Python · Lambda with ReportBatchItemFailuresretry only failures
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.

A SLOW BATCH, A REBALANCE, AND A REPLAY 0 s0–300 s300 s301 s320 s Consumer A Consumer A Coordinator Consumer B Consumer A poll() → offsets 100–599 · committed = 100 processing slowly: a downstream API is degraded max.poll.interval.ms exceeded → rebalance assigned partition 3 → resumes at offset 100 commitSync() → CommitFailedException Result: offsets 100 onward processed by A and by B A’s work was never committed, so B redid it, partly while A was still running.
A slow dependency turns into duplicate processing: the batch outlives 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.

FeatureWhat it guaranteesWhat it does not cover
SQS FIFO deduplicationA second send with the same deduplication id within 5 minutes is accepted but not deliveredRedelivery after a visibility timeout or crash; resends after the 5-minute interval
Kafka idempotent producerProducer retries do not write duplicate records to a partitionThe application publishing the same event twice
Kafka transactions / Streams exactly_once_v2Read-process-write to Kafka topics commits output and offsets atomically, for readers using read_committedWrites 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:

Python · SQS consumerduplicates charge twice
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:

  1. 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.
  2. Record the id under a unique constraint, in the same database transaction as the consumer’s state change.
  3. Guard every state transition so it only applies from the expected previous state.
  4. Make external calls idempotent with a key derived from the business event, or hand them to an outbox worker that does.
  5. Acknowledge only after commit. Delete the message or commit the offset last.
AN IDEMPOTENT CONSUMER, STEP BY STEP 1 Receive; read event_id from the payload 2 BEGIN; INSERT INTO processed_messages ON CONFLICT DO NOTHING RETURNING event_id 3 Guarded state change + outbox row UPDATE … WHERE status = 'pending' 4 COMMIT 5 Acknowledge: DeleteMessage / commit offset 0 ROWS Duplicate delivery COMMIT, acknowledge, stop Crash before step 4 Everything rolls back; redelivery starts clean Crash between 4 and 5 Redelivery hits the dedupe row and is just acknowledged Payments and emails run from the outbox with their own idempotency keys, outside this transaction.
The dedupe record and the state change commit together. A redelivery at any point either starts clean or is recognized and acknowledged.
Python · idempotent handlersafe to run twice
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:

Java · Kafka consumercommit after work
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.

KeyCatchesMisses
SQS MessageIdRedelivery of the same sent messageProducer resends, which get a new MessageId
Kafka topic + partition + offsetReplays after a crash or rebalanceThe same event published twice, which gets a new offset
Hash of the payloadByte-identical copiesCopies with a new timestamp or trace id; and it wrongly merges two genuine identical events
Producer-assigned event idRedeliveries, replays, and resendsOnly 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.

SQS
  • Visibility timeout above worst-case processing for a whole batch
  • Heartbeat with ChangeMessageVisibility for long jobs
  • Lambda: ReportBatchItemFailures, and a queue visibility timeout at least 6× the function timeout
  • Graceful shutdown drains in-flight messages
Kafka
  • Manual commit after processing
  • Smaller max.poll.records so a batch fits inside max.poll.interval.ms
  • Commit in onPartitionsRevoked of a ConsumerRebalanceListener
  • Timeouts on downstream calls so one slow call cannot stall a batch
Both
  • 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
Longer timeouts have a cost

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:

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:

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

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 →