Duplicate Events and Idempotent Consumers: At-Least-Once Done Right
At-least-once delivery means duplicates. Dedupe with a stable key in the same transaction, ack after commit, and guard external side effects.
Azeem Subhani · · 13 min read

A customer was charged twice, or a seat was reserved twice, and the logs show the same event id processed by your consumer two times. The first instinct is to file a ticket against the queue, the webhook provider, or the stream: the bus delivered a message twice, so the bus is broken. It is not. Almost every queue, webhook system, and stream you are likely to use promises at-least-once delivery, and a duplicate is that promise working as written. Duplicate events are the input your consumer must handle, and an idempotent consumer is the only fix that holds up. Below: where the second delivery comes from, why broker identifiers make poor dedup keys, how to apply each event once, and the rule that tells you which side effects can still happen twice.
Why duplicate events are not a bus bug
Read the delivery documentation of the systems you already run and the duplicate is in writing:
- Amazon SQS standard queues store copies of each message on multiple servers. If one of those servers is unavailable when you delete a message, that copy survives and can be delivered again. AWS tells you to design applications to be idempotent (SQS at-least-once delivery).
- Stripe says webhook endpoints "might occasionally receive the same event more than once" and tells you to log the event ids you have processed and skip ones you have seen (Stripe webhooks).
- Google Cloud Pub/Sub, even with exactly-once delivery turned on, says a subscription can still receive multiple copies of a message because of publish-side duplicates (Pub/Sub exactly-once delivery).
You do not need a broker fault to get a duplicate, either. A common cause is your own consumer. It receives a message, performs the side effect, and then crashes, times out, or loses its connection before it acknowledges. The broker never heard "done," so it delivers again. In SQS terms, the message becomes visible again when the visibility timeout expires before you delete it, and the default timeout for a queue is 30 seconds (SQS visibility timeout). A slow downstream call on a busy day is enough.
Retries stack on top of that. A producer that retries after a lost acknowledgement can publish the same logical event twice. AWS lists retrying non-idempotent calls, and the duplicated results that follow, as a named anti-pattern (AWS Well-Architected REL05-BP03). If retries are also firing at several layers at once, read retry storms, backoff, and jitter before you tune anything here.
At least once, at most once, and what exactly once would take
The three delivery terms are about where the consumer records progress relative to the work.
- At most once. Record progress, then do the work. A crash after recording and before the work loses the message. Confluent's description of Kafka consumer semantics defines it by that order: save the position first, process second (Kafka delivery semantics).
- At least once. Do the work, then record progress. A crash between the two repeats the work. Same source, opposite order: process first, save the position second.
- Exactly once. The work and the progress marker commit together, atomically, or neither does.
The third one is not a broker setting. It is a property of the pair: the system that holds your progress and the system that holds your side effect. When both live in one transactional store, you can get it. Kafka does this when both the consumer position and the output are Kafka topics: the offset commit and the output records go into one Kafka transaction. Kafka Connect gets the same effect for HDFS by storing offsets next to the data it writes, so the two update together or not at all (Kafka delivery semantics).
The guarantee stops where that shared transaction stops. Databricks states this boundary plainly for its pipelines: exactly-once applies to managed Delta-to-Delta flows, and for foreach_batch_sink and custom external writes you should treat the edge as at-least-once, because a batch that is retried after a partial write can leave duplicate rows in the external system. Their fix is to make the external write idempotent, either by upserting on a natural key or by writing a batch id the receiver can dedupe on (Databricks processing guarantees).
SQS FIFO is a useful second example of a scoped promise. FIFO queues don't introduce duplicates if you retry SendMessage within the 5-minute deduplication interval (SQS exactly-once processing). That protects the queue from a retrying producer. It says nothing about a consumer that charges a card and then crashes before deleting the message.
So stop reading "exactly once" as a checkbox. Ask whether one transaction covers both the progress marker and the effect.
Why an offset or message id is a poor dedup key
The natural move is to dedupe on whatever identifier the broker hands you. Most of them identify a delivery or a position, not the business event.
- SQS receipt handles change on every receive. AWS documents that a message received twice comes with a different receipt handle each time (SQS message identifiers).
- SQS message ids are assigned per
SendMessagecall. A producer that retries a send creates a second message with a second id. - Pub/Sub publish-side duplicates come from publish retries, by the client or by the service itself, and they arrive as separate messages (Pub/Sub exactly-once delivery).
- Kafka offsets identify a position in one partition of one topic. If the producer wrote the event twice, it sits at two offsets. If you replay into a new topic, rebuild a topic, or move consumers to a mirrored cluster, the offsets change while the event does not.
- Stripe warns that it sometimes generates two separate Event objects for one change, and suggests identifying those duplicates by the id of the object in
data.objecttogether withevent.type. It also says not to use thecreatedtimestamp to decide whether you have already processed an event (Stripe webhooks).
The dedup key should be the identity of the change in your domain: a producer-assigned event id that survives retries, an order id plus an operation name, or a provider object id plus event type. Ask the producer to mint the id once, before the first send attempt, and to reuse it on every retry.
How to confirm you are seeing redelivery
Before you change code, prove which mechanism produced the duplicate. The evidence tells you how long your dedup window must be.
- Pull both processing records for the duplicated event. Log the domain event id, the broker message id, the receive attempt count if your broker exposes one, and the timestamps of receive, side effect, and acknowledgement. Compare the two records.
- Same message id, different receipt handle or attempt count: classic redelivery. Look at the first attempt. Did it exceed the visibility timeout or ack deadline, crash, or fail to acknowledge?
- Different message ids, same domain id: the producer published twice. Check producer retries and any outbox relay that might resend after a crash.
- Two provider events for the same object and type (common with webhooks): the provider generated two events. Your key must be the object id plus type, not the event id alone.
- Measure the gap between the two attempts. Hours or days apart usually means a dead-letter redrive or a manual resend, not an automatic retry. That gap sets the minimum retention for your dedup records.
- Check where the acknowledgement happens in the code path. If the ack runs before the state change commits, you also have an at-most-once bug hiding behind the duplicate: a crash at the wrong moment would lose the event instead.
Building an idempotent consumer for duplicate events
The pattern has three parts. Record the dedup key in the same database transaction as the state change. Acknowledge the message only after that transaction commits. Keep the dedup records longer than any path that can redeliver.
Start with a table that holds one row per applied event per consumer. Scoping by consumer matters: two different consumers that read the same event each need to apply it once.
-- Illustrative schema. One row per (consumer, event) that has been applied.
CREATE TABLE processed_events (
consumer text NOT NULL,
event_key text NOT NULL,
processed_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (consumer, event_key)
);
CREATE INDEX processed_events_processed_at_idx ON processed_events (processed_at);
The consumer inserts the dedup row and applies the change in one transaction. If the insert conflicts, the event was already applied, so the consumer rolls back and acknowledges. The example uses pg and the AWS SDK v3 SQS client.
// Illustrative consumer: Postgres state change + dedup row in one transaction,
// acknowledge (DeleteMessage) only after COMMIT succeeds.
import { Pool } from "pg";
import { SQSClient, ReceiveMessageCommand, DeleteMessageCommand } from "@aws-sdk/client-sqs";
const pool = new Pool();
const sqs = new SQSClient({});
const QUEUE_URL = process.env.RESERVATIONS_QUEUE_URL!;
type SeatHeld = { eventId: string; seatId: string; bookingId: string };
async function applyOnce(evt: SeatHeld): Promise<"applied" | "duplicate"> {
const db = await pool.connect();
try {
await db.query("BEGIN");
const dedup = await db.query(
`INSERT INTO processed_events (consumer, event_key)
VALUES ('reservations', $1)
ON CONFLICT (consumer, event_key) DO NOTHING
RETURNING event_key`,
[evt.eventId],
);
if (dedup.rowCount === 0) {
// Already applied by an earlier delivery. Nothing to do.
await db.query("ROLLBACK");
return "duplicate";
}
const held = await db.query(
`UPDATE seats SET booking_id = $2
WHERE seat_id = $1 AND booking_id IS NULL`,
[evt.seatId, evt.bookingId],
);
if (held.rowCount === 0) {
// Seat taken by someone else: a business outcome, not a retryable error.
// Record it so redelivery does not re-run this decision.
await db.query(
`INSERT INTO booking_failures (booking_id, reason) VALUES ($1, 'seat_taken')`,
[evt.bookingId],
);
}
await db.query("COMMIT");
return "applied";
} catch (err) {
await db.query("ROLLBACK").catch(() => {});
throw err; // Do not ack. The broker will redeliver after the visibility timeout.
} finally {
db.release();
}
}
export async function pollOnce() {
const res = await sqs.send(new ReceiveMessageCommand({
QueueUrl: QUEUE_URL, MaxNumberOfMessages: 10, WaitTimeSeconds: 20,
}));
for (const msg of res.Messages ?? []) {
const evt = JSON.parse(msg.Body!) as SeatHeld;
const outcome = await applyOnce(evt);
console.log(JSON.stringify({ eventId: evt.eventId, messageId: msg.MessageId, outcome }));
// Ack after commit, for "applied" and "duplicate" alike.
await sqs.send(new DeleteMessageCommand({
QueueUrl: QUEUE_URL, ReceiptHandle: msg.ReceiptHandle!,
}));
}
}
A few details carry the correctness:
- Concurrent duplicates. If two workers receive the same event at once, both try to insert the same primary key. The second insert waits on the first transaction's uncommitted row. If the first commits, the second sees the conflict and does nothing. If the first rolls back, the second proceeds. You get one application without a separate lock.
- Ack after commit. A crash between
COMMITandDeleteMessageproduces a redelivery, which the dedup row absorbs. A crash beforeCOMMITrolls back both the dedup row and the state change, so the redelivery applies the event cleanly. Neither order of failure loses or doubles the reservation. - Business failures are outcomes. "Seat already taken" commits with the dedup row. If you threw instead, the message would bounce until it reached the dead-letter queue, and a later redrive could apply it in a different world.
- Duplicate is not an error. Log it as an outcome and chart it. A sudden rise points at a timeout or a misbehaving producer.
Prefer a versioned upsert for state updates
When the event carries the full new state of an entity, you often don't need a dedup table at all. Write the state with a version and let the database discard anything older. PostgreSQL documents that ON CONFLICT DO UPDATE always ends in exactly one of insert or update, even under concurrency, and the optional WHERE clause means rows that don't satisfy it are locked but not updated (PostgreSQL INSERT).
-- Illustrative: apply a "subscription state" event only if it is newer.
INSERT INTO subscription_state AS s (subscription_id, status, version, updated_at)
VALUES ($1, $2, $3, now())
ON CONFLICT (subscription_id) DO UPDATE
SET status = EXCLUDED.status,
version = EXCLUDED.version,
updated_at = now()
WHERE s.version < EXCLUDED.version;
This handles both duplicates and out-of-order delivery, which matters because providers such as Stripe do not guarantee ordering. It only works when the producer supplies a version that increases monotonically per entity. A wall-clock timestamp is not that: Stripe records created in seconds, so distinct events can share a timestamp. If you have no reliable version, a common alternative for webhooks is to treat the event as a hint and fetch the current object from the provider's API before writing it.
Webhooks: acknowledge fast, but only after a durable write
Stripe tells you to return a 2xx quickly, before any complex logic that might time out (Stripe webhooks). That looks like it conflicts with "ack after commit." It doesn't, as long as the thing you commit before acknowledging is small: verify the signature, insert the event into an inbox table keyed by the dedup key, commit, return 2xx. A separate worker then processes the inbox with the transaction pattern above. The acknowledgement now means "durably received," and the inbox insert absorbs duplicate deliveries. The durable webhook acknowledgement post covers the inbox in detail.
The side effect outside the transaction
Here is the rule that decides whether your consumer is actually safe: name every side effect that is not in the same transaction as the dedup row. Each of those can still happen twice.
Apply it to two effects in one booking flow.
The reservation is safe. The seat update and the dedup row commit together in Postgres. Either both exist or neither does. A redelivery finds the dedup row and stops.
The charge is not. Calling a payment provider is a network request to a system your transaction does not cover. Consider the sequence: insert dedup row, call the provider, the provider charges the card, your process crashes before COMMIT. The transaction rolls back, the dedup row disappears, the message is redelivered, and the consumer charges again. Moving the provider call after COMMIT doesn't fix it either: a crash between the commit and the call means the redelivery finds the dedup row, stops, and the charge never happens. Holding a database transaction open across that remote call has its own costs, covered in long transactions around external calls.
The fix is to push idempotency into the external call and keep a local record of intent.
- In one transaction, insert a
payment_attemptsrow keyed by the business operation (for example, the order id pluscharge), with statuspendingand an idempotency key derived from that operation, not from the delivery. For this consumer, that row is the dedup record. - Commit, then call the provider with that key. With Stripe's Node library, that is
stripe.paymentIntents.create(params, { idempotencyKey }). - Record the provider's result on the
payment_attemptsrow, then acknowledge the message. - On redelivery, read the
payment_attemptsrow instead of stopping at "already seen." If it issucceededorfailed, acknowledge and stop. If it is stillpending, a previous attempt died somewhere around the call: reconcile against the provider before deciding to call again. - Run a periodic sweep for rows stuck in
pending, because the message that would have resumed them may have gone to the dead-letter queue.
Step 4 matters because the provider's idempotency has a time limit. Stripe saves the status code and body of the first request for a key and replays it on retries, but keys can be removed once they are at least 24 hours old, and a key reused after pruning is treated as a brand-new request (Stripe idempotent requests). Stripe itself retries failed webhook deliveries for up to three days in live mode, and a manual resend works up to 15 days (Dashboard) or 30 days (CLI) after the event was created (Stripe webhooks). A redelivery on day five with the same idempotency key can therefore create a new charge. Your local intent row is the only record that outlives that window. For the race inside the idempotency layer itself, see idempotency key race conditions.
Emails, SMS, and third-party calls without idempotency keys fall in the same bucket. For most notifications a rare duplicate is acceptable; decide that on purpose.
Trade-offs and when not to use each fix
At-least-once plus an idempotent apply loses nothing. It costs a dedup table, a transaction that includes it, and discipline about where the ack goes. It is the right default for anything that moves money, inventory, or user-visible state.
At-most-once (ack first, then work) is simpler and drops work on every crash. Use it only for telemetry, cache warming, or other work where a loss is cheaper than a duplicate and nobody will reconcile.
A broker's exactly-once setting helps inside its boundary. FIFO deduplication protects against producer retries within its interval. Kafka transactions cover topic-to-topic processing. Pub/Sub exactly-once applies only to pull subscriptions, and only when subscribers connect in the same region (Pub/Sub exactly-once delivery). None of them covers a charge, an email, or a row in a database outside that transaction. Turning them on is fine; relying on them for external effects is not.
The dedup table grows. It needs a retention rule, and the rule has a failure mode: delete a key too early and a late redelivery is applied again. Set retention longer than the longest redelivery path you have: automatic retries, dead-letter retention plus the time it takes someone to redrive, and any manual resend window your provider offers. Prune by processed_at in small batches so the delete does not hold locks on a hot table.
The versioned upsert is the cheapest option and needs no extra table, but only works for full-state events with a trustworthy version. It is wrong for deltas such as "add 5 credits," where applying the event twice changes the result.
Inbox tables for webhooks decouple acknowledgement from processing, but an inbox that stops draining is an outage that looks like silence. Alert on inbox age, not only on errors.
Checklist for duplicate events
- For each consumer, write down the dedup key and confirm it is a domain identity that survives producer retries, not a receipt handle, offset, or delivery id.
- Confirm the dedup row and the state change commit in one transaction.
- Move the acknowledgement after
COMMIT. Ack duplicates too. - List every side effect outside that transaction: payment calls, emails, third-party writes. For each, either pass a business-derived idempotency key and keep a local intent row, or accept and document the duplicate.
- Set dedup retention longer than the longest redelivery window, including dead-letter redrives and manual resends.
- Log
appliedandduplicateas outcomes and chart the duplicate rate. - Test it: deliver the same message twice in a row, deliver it twice concurrently, and kill the consumer between the side effect and the ack. The state should be identical in all three cases.
Sources
- Amazon SQS at-least-once delivery
- Amazon SQS visibility timeout
- Amazon SQS queue and message identifiers
- Exactly-once processing in Amazon SQS
- AWS Well-Architected, REL05-BP03 Control and limit retry calls
- Google Cloud Pub/Sub, exactly-once delivery
- Confluent, Kafka message delivery guarantees
- Databricks, processing guarantees in Lakeflow pipelines
- Stripe, receive Stripe events in your webhook endpoint
- Stripe API, idempotent requests
- PostgreSQL, INSERT and ON CONFLICT
Written by
Azeem Subhani
Senior Full-Stack & AI Application Engineer
I build SaaS, booking, payment, real-time, and AI-enabled web platforms with React, Next.js, Node.js, NestJS, Django, PostgreSQL, and AWS. My work includes Stripe payment systems, white-label booking flows, real-time collaboration, RAG workflows, and developer automation.


