Backend Systems

Idempotent Pub/Sub consumers: exactly-once is your job

Pub/Sub will deliver some messages twice, and that is by design. Here is how to build consumers whose side effects happen once anyway, with working code.

Published

Updated

—

Reading time

17 min

Every team that moves side effects onto a message queue meets the same bug sooner or later. A counter ends up at 103 when it should be 101. A customer gets the same receipt email twice. A webhook fires three times and a partner's system books three orders. Nobody wrote a loop; the broker delivered the same message more than once, which it is allowed to do, and the consumer treated each delivery as new. This article is for engineers running consumers on Google Cloud Pub/Sub, though almost everything here applies to SQS, Kafka consumers, or any at-least-once queue. I'll cover why duplicates happen, what Pub/Sub's exactly-once feature does and doesn't give you, and how I build consumers so that duplicate deliveries have no visible effect.

The constraint: at-least-once is the contract#

Pub/Sub's default delivery guarantee is at-least-once. Each published message is delivered at least once to each subscription. It may be delivered more than once, and your code has to handle that. Duplicates come from several independent places, and closing one of them leaves the others open.

Ack deadline expiry. Every delivery starts a timer, the ack deadline. If your consumer doesn't acknowledge before it runs out, Pub/Sub assumes the message was lost and sends it again. A slow downstream call, a garbage-collection pause, or a cold start can push you past the deadline while the first attempt is still running. You then have two copies in flight at once.

Explicit retries. When a push endpoint returns a non-success status, or a pull client nacks, the message goes back into the queue. If the handler had already done half its work before failing (updated the database but then timed out calling the email provider), the retry repeats that half.

Crash after effect, before ack. The process writes to the database, then the container is killed during a scale-in or deploy before the ack goes out. From Pub/Sub's point of view nothing was acknowledged, so the message is redelivered. From your database's point of view, the work is already done.

Publisher retries. This one surprises people. If a publish call times out on the client even though the server received it, the client library retries and the topic can end up holding two messages with the same payload and different message IDs. No consumer-side feature that dedupes by message ID can catch it, because as far as Pub/Sub is concerned they are two different messages.

Subscription-level replays. Seek to a timestamp or snapshot redelivers everything after that point. That's intended behaviour, and your consumers have to tolerate it too.

The conclusion I've come to is simple: the queue delivers messages, and whether the effect happens once is my code's responsibility.

What Pub/Sub's exactly-once delivery actually covers#

Pub/Sub has an exactly-once delivery option on subscriptions. It's useful, and it's also narrower than the name suggests. As of September 2026, the documented scope is:

  • Pull subscriptions only, including StreamingPull. Push and export subscriptions (BigQuery, Cloud Storage) don't support it.
  • Regional. The guarantee holds within a single cloud region. If your subscribers connect from several regions, a message can still be delivered to more than one of them.
  • Keyed on message ID. Once a message has been acknowledged successfully, it won't be redelivered, and it won't be redelivered while an ack deadline is still outstanding. Duplicates created by publisher retries have different message IDs, so they are delivered normally.
  • Acks can fail, and you have to check. With exactly-once enabled, an ack is a round trip that returns success or failure. In the Node.js client that means calling message.ackWithResponse() and handling the rejected promise. If the ack fails because the deadline expired, the message will come back, and your side effect may already have run.
  • It costs latency. The documentation is explicit that publish-to-subscribe latency is significantly higher with exactly-once than without it.

None of this makes your side effects exactly-once. Pub/Sub can promise not to hand you the same message ID again after a successful ack, but it can't know whether your database commit happened before the process died. It's a good way to reduce the duplicate rate. Correctness still has to come from your consumer.

NoteEnd-to-end is the only exactly-once that matters

"Exactly-once delivery" describes the broker. What users see is exactly-once effect: one counter increment, one email, one charge. You get that from an idempotent consumer, whatever the broker promises.

Options for making the effect happen once#

There are four main approaches. They aren't mutually exclusive, and most real consumers combine two of them.

ApproachCatchesCost / complexityBest for
Naturally idempotent writes (set, upsert by key)All duplicate sourcesNone; it's a modelling choiceState you can express as "make it so"
Dedupe record in the same transaction as the effectAll duplicates that share the idempotency keyOne extra read + write per messageEffects in the same database as the dedupe store
Dedupe store beside the effect (Redis SET NX, separate table)Most duplicates; a gap exists between effect and markerExtra network hop; TTL tuningHigh-volume, low-stakes effects
Claim/lease + downstream idempotency key (outbox-style)All duplicates, if the downstream honours keysMost moving partsExternal calls: email, payments, webhooks
Pub/Sub exactly-once deliveryRedeliveries of the same message ID onlyHigher latency; pull only; regionalReducing duplicate rate, never as the only guard

Choosing the idempotency key#

The key matters more than where you store it. There are two families.

An event ID is generated once, at the moment the fact happens, and carried in the payload (a UUIDv7 works well). It's stable across publisher retries only if you generate it before the publish call and outside any retry loop. Pub/Sub's own messageId does not qualify, because a publisher retry produces a new one.

A natural business key is derived from the domain: order:8812:receipt-email, like:{postId}:{userId}, invoice:2026-09:{customerId}. It catches duplicates that an event ID can't, such as a user double-clicking and producing two distinct events for the same intent, or two services emitting the same fact independently.

My rule is to use the natural key whenever the domain has one ("this customer gets one receipt per order"), and the event ID when the event itself is the unit ("apply this +1"). In both cases, prefix the key with the consumer name. Two consumers on two subscriptions of the same topic each need their own record of what they've processed.

The decision#

For effects that live in the same database as my dedupe records, I put the dedupe record in the same transaction as the effect. That way the marker and the effect commit together or not at all, and the ordering bugs can't happen. For effects that leave the building (email, third-party APIs, payments), I use a lease-based claim plus the downstream provider's idempotency key where one exists. I turn on Pub/Sub exactly-once only for pull consumers where a duplicate is expensive enough to justify the extra latency, and I treat it as a way to lower the duplicate rate, never as the thing that makes the consumer correct.

Push consumer with a transactional dedupe guard, backoff retries and a dead-letter topic.

Implementation#

The example is a push consumer on Cloud Run written with Fastify and the Firestore Node SDK. It keeps a denormalised likeCount on a post document. The API writes the like itself synchronously and publishes a like.added event; the consumer does the counter update asynchronously.

Step 1: put a stable ID on the event at publish time#

src/events/publish.tsts
import { PubSub } from "@google-cloud/pubsub";
import { randomUUID } from "node:crypto";
 
const pubsub = new PubSub();
const topic = pubsub.topic("like-events");
 
export async function publishLikeAdded(postId: string, userId: string) {
  const event = {
    // Generated once, before any retry, so publisher retries carry the same ID.
    eventId: randomUUID(),
    type: "like.added",
    postId,
    userId,
    occurredAt: new Date().toISOString(),
  };
  // The client library retries transient failures internally; eventId stays stable.
  await topic.publishMessage({ json: event, attributes: { type: event.type } });
}

If publishing must be tied atomically to a database write, use a transactional outbox: write the event row in the same transaction as the business change, and have a relay publish it. The event ID is then the outbox row ID, which is stable by construction.

Step 2: the push handler and its ack semantics#

For push subscriptions, the HTTP status is the ack. A 102, 200, 201, 202, or 204 response acknowledges the message; any other status, or no response before the ack deadline, counts as a nack and triggers a retry. That gives the handler three possible outcomes:

  • Success or duplicate: return 204.
  • Transient failure (database contention, a downstream timeout): return 503 so Pub/Sub retries with backoff.
  • Poison message (unparseable, fails validation, references data that will never exist): retrying won't help. Put it somewhere a human can look at it, then return 204.
src/server.tsts
import Fastify from "fastify";
import { Firestore, FieldValue, Timestamp } from "@google-cloud/firestore";
 
const db = new Firestore();
const app = Fastify({ logger: true });
 
type PushBody = {
  message: {
    data?: string; // base64
    attributes?: Record<string, string>;
    messageId: string;
    publishTime: string;
    orderingKey?: string;
  };
  subscription: string;
  deliveryAttempt?: number; // present when a dead-letter policy is set
};
 
type LikeAdded = { eventId: string; type: "like.added"; postId: string; userId: string };
 
class PoisonMessage extends Error {}
 
function parseLikeAdded(message: PushBody["message"]): LikeAdded {
  if (!message.data) throw new PoisonMessage("empty data");
  let body: unknown;
  try {
    body = JSON.parse(Buffer.from(message.data, "base64").toString("utf8"));
  } catch {
    throw new PoisonMessage("invalid JSON");
  }
  const e = body as Partial<LikeAdded>;
  if (e.type !== "like.added" || !e.eventId || !e.postId || !e.userId) {
    throw new PoisonMessage("schema mismatch");
  }
  return e as LikeAdded;
}
 
const CONSUMER = "like-counter";
const DEDUPE_TTL_MS = 7 * 24 * 60 * 60 * 1000; // must exceed the subscription's retention window
 
async function applyOnce(evt: LikeAdded): Promise<"applied" | "duplicate" | "skipped"> {
  const markerRef = db.collection("processedEvents").doc(`${CONSUMER}:${evt.eventId}`);
  const postRef = db.collection("posts").doc(evt.postId);
 
  return db.runTransaction(async (tx) => {
    // Firestore transactions require all reads before any writes.
    const [marker, post] = await Promise.all([tx.get(markerRef), tx.get(postRef)]);
    if (marker.exists) return "duplicate";
 
    // Record the event even when the post is gone, so a redelivery is still a no-op.
    tx.create(markerRef, {
      processedAt: FieldValue.serverTimestamp(),
      expireAt: Timestamp.fromMillis(Date.now() + DEDUPE_TTL_MS),
    });
    if (!post.exists) return "skipped";
 
    tx.update(postRef, { likeCount: FieldValue.increment(1) });
    return "applied";
  });
}
 
app.post<{ Body: PushBody }>("/pubsub/like-events", async (req, reply) => {
  const { message, deliveryAttempt, subscription } = req.body;
  const log = req.log.child({ messageId: message.messageId, deliveryAttempt, subscription });
 
  let evt: LikeAdded;
  try {
    evt = parseLikeAdded(message);
  } catch (err) {
    // Poison: quarantine for inspection, then ack so it stops consuming retries.
    await db.collection("quarantine").doc(`${CONSUMER}:${message.messageId}`).set({
      reason: (err as Error).message,
      message,
      receivedAt: FieldValue.serverTimestamp(),
    });
    log.error({ err }, "poison message quarantined");
    return reply.code(204).send();
  }
 
  try {
    const outcome = await applyOnce(evt);
    log.info({ eventId: evt.eventId, outcome }, "like event handled");
    return reply.code(204).send();
  } catch (err) {
    log.warn({ err, eventId: evt.eventId }, "transient failure; asking Pub/Sub to retry");
    return reply.code(503).send();
  }
});
 
app.listen({ host: "0.0.0.0", port: Number(process.env.PORT ?? 8080) });

The highlighted function is where the guarantee comes from. The read of the marker, the creation of the marker, and the increment all commit atomically. If two deliveries of the same event run concurrently, Firestore's optimistic concurrency makes one transaction retry; on the retry it sees the marker and returns "duplicate". If the process dies before commit, nothing was written, and the redelivery does the work from scratch.

Set a Firestore TTL policy on processedEvents.expireAt so the dedupe collection doesn't grow forever. TTL deletion isn't immediate (the documentation says expired documents are typically removed within about a day), so don't write logic that depends on a marker disappearing at a precise time. Make sure the TTL is comfortably longer than the subscription's message retention duration plus any replay window you might use. If a marker expires while its message can still be redelivered, the guard has a hole.

On Cloud Run, you don't need to verify the push request's OIDC token in code. Deploy the service with --no-allow-unauthenticated, give the subscription a push service account, and grant that account roles/run.invoker on the service. Cloud Run's front end rejects unauthenticated requests before your handler runs.

Step 3: the subscription, retry policy and dead-letter topic#

infra/subscription.shbash
PROJECT_ID="my-project"
PROJECT_NUMBER="$(gcloud projects describe "$PROJECT_ID" --format='value(projectNumber)')"
PUBSUB_SA="service-${PROJECT_NUMBER}@gcp-sa-pubsub.iam.gserviceaccount.com"
PUSH_SA="pubsub-push@${PROJECT_ID}.iam.gserviceaccount.com"
URL="$(gcloud run services describe like-consumer --region=europe-west1 --format='value(status.url)')"
 
# Dead-letter topic, plus a pull subscription so dead letters are retained and inspectable.
gcloud pubsub topics create like-events-dlq
gcloud pubsub subscriptions create like-events-dlq-inspect \
  --topic=like-events-dlq --message-retention-duration=7d
 
gcloud pubsub subscriptions create like-events-counter \
  --topic=like-events \
  --push-endpoint="${URL}/pubsub/like-events" \
  --push-auth-service-account="$PUSH_SA" \
  --ack-deadline=60 \
  --min-retry-delay=10s \
  --max-retry-delay=600s \
  --dead-letter-topic=like-events-dlq \
  --max-delivery-attempts=10
 
# The Pub/Sub service agent must be able to publish to the DLQ and ack on the source subscription.
gcloud pubsub topics add-iam-policy-binding like-events-dlq \
  --member="serviceAccount:${PUBSUB_SA}" --role=roles/pubsub.publisher
gcloud pubsub subscriptions add-iam-policy-binding like-events-counter \
  --member="serviceAccount:${PUBSUB_SA}" --role=roles/pubsub.subscriber
 
gcloud run services add-iam-policy-binding like-consumer --region=europe-west1 \
  --member="serviceAccount:${PUSH_SA}" --role=roles/run.invoker

Some notes on these flags:

  • Ack deadline. For push, this is how long Pub/Sub waits for your HTTP response. Set it above your handler's worst realistic latency (cold start included), but not so high that a hung request holds a message for minutes. Keep your service's request timeout at or below it, so a slow request fails visibly instead of producing a duplicate in the background.
  • Retry policy. Exponential backoff between the minimum and maximum delay. Without it, Pub/Sub redelivers as soon as it can, which during a downstream outage becomes a retry storm against a system that's already struggling.
  • Dead-letter topic. max-delivery-attempts accepts 5 to 100 (default 5). The attempt count is best-effort, so treat it as approximate. A dead-letter topic with no subscription drops messages, so always attach one and alert on its backlog (subscription/num_undelivered_messages in Cloud Monitoring).
CostWhat the guard costs

The transactional guard adds one document read and one document write per message, plus the TTL delete later. At the time of writing, Firestore bills reads, writes and deletes per operation, so a consumer handling 10 million messages a month pays for about 10 million extra reads and writes. That is usually small next to what a duplicate costs (a double charge, a support ticket, a corrupted metric), but check it against your volume. For very high-throughput, low-stakes counters, a Redis SET NX guard or batched aggregation may be cheaper.

Step 4: the other dedupe stores#

If the effect isn't in Firestore, the pattern stays the same and only the store changes.

Relational database, same transaction. A unique constraint does the dedupe work:

migrations/001_processed_events.sqlsql
CREATE TABLE processed_events (
  consumer     text        NOT NULL,
  event_id     text        NOT NULL,
  processed_at timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (consumer, event_id)
);
src/consumers/pg.tsts
import type { PoolClient } from "pg";
 
export async function applyOncePg(client: PoolClient, consumer: string, eventId: string, effect: () => Promise<void>) {
  await client.query("BEGIN");
  try {
    const res = await client.query(
      "INSERT INTO processed_events (consumer, event_id) VALUES ($1, $2) ON CONFLICT DO NOTHING",
      [consumer, eventId],
    );
    if (res.rowCount === 0) {
      await client.query("ROLLBACK");
      return "duplicate" as const;
    }
    await effect(); // must use the same client so it joins this transaction
    await client.query("COMMIT");
    return "applied" as const;
  } catch (err) {
    await client.query("ROLLBACK");
    throw err;
  }
}

Redis SET NX with a TTL. This is fast and cheap, but it isn't transactional with your effect:

src/consumers/redis-guard.tsts
import Redis from "ioredis";
const redis = new Redis(process.env.REDIS_URL!);
 
// Returns true if this caller won the key; false if already seen.
export async function firstSeen(key: string, ttlSeconds: number): Promise<boolean> {
  return (await redis.set(key, "1", "EX", ttlSeconds, "NX")) === "OK";
}

The weakness is the ordering. If you set the key first and then crash before the effect, the retry sees the key and skips, and the effect is lost. If you run the effect first and set the key afterwards, a crash between the two produces a duplicate. Pick whichever failure you can tolerate, and use this only for effects where that trade-off is acceptable: cache warming, analytics fan-out, notifications you'd rather drop than repeat.

Step 5: side effects that leave your database#

An email or a third-party API call can't join your database transaction. For those I use a claim with a lease, plus the provider's idempotency key where it has one (payment APIs such as Stripe accept an Idempotency-Key header; check what your email or webhook provider supports).

src/consumers/effect-lease.tsts
import { Firestore, FieldValue } from "@google-cloud/firestore";
const db = new Firestore();
 
type Claim = "claimed" | "done" | "busy";
 
export async function claim(key: string, leaseMs: number): Promise<Claim> {
  const ref = db.collection("effects").doc(key);
  return db.runTransaction(async (tx) => {
    const snap = await tx.get(ref);
    const now = Date.now();
    if (!snap.exists) {
      tx.create(ref, { status: "pending", leaseUntil: now + leaseMs, attempts: 1 });
      return "claimed";
    }
    const d = snap.data()!;
    if (d.status === "done") return "done";
    if (d.leaseUntil > now) return "busy"; // someone else is working on it right now
    tx.update(ref, { leaseUntil: now + leaseMs, attempts: FieldValue.increment(1) });
    return "claimed"; // previous holder's lease expired; take over
  });
}
 
export async function markDone(key: string) {
  await db.collection("effects").doc(key).update({ status: "done", doneAt: FieldValue.serverTimestamp() });
}

The handler calls claim("receipt-email:order:8812", 30_000). On "done" it returns 204. On "busy" it returns 503, so Pub/Sub backs off and tries again later. On "claimed" it sends the email with the same key as the provider's idempotency key, calls markDone, and returns 204. One window remains: if the provider accepts the send and the process dies before markDone, the next holder sends again. The provider-side key closes that window. Without one, the window is small but not zero, and you should say so in the design doc rather than claim exactly-once.

Counters have a simpler fix than leases: turn "increment" into "set membership". Instead of likeCount += 1, the source of truth is a document per (postId, userId), created with a deterministic ID. A duplicate create is a no-op by construction. The count is then either derived (with an aggregation query) or maintained by the guarded increment above.

Trade-offs and failure modes#

Ordering keys don't remove duplicates, and they can add some. With message ordering enabled, messages with the same ordering key are delivered in publish order within a region. When one message in a key's sequence is nacked or its deadline expires, Pub/Sub redelivers it and the messages after it with the same key, so a single slow handler can produce a run of duplicates. A poison message also blocks every later message on its key until it's acknowledged or dead-lettered, and once it's dead-lettered, ordering continues past the gap. If your consumer relies on order for correctness, carry a version or sequence number in the payload and reject stale writes (if incoming.version <= stored.version, skip). That's idempotent and order-tolerant at once, and it protects you from producers that don't use ordering keys correctly.

Contention on hot keys. Firestore transactions on the same document serialise. A viral post receiving thousands of like events per second will cause transaction retries and latency, because the dedupe guard adds a second document to the contended transaction. The standard fixes are distributed counters (sharded counter documents), or aggregating in the consumer and flushing a batch of deltas under a single dedupe key per batch.

Treating every failure as transient. If the handler returns 503 for a validation error, that message uses up all its delivery attempts, each one with backoff, before it reaches the dead-letter topic, and it adds noise to your retry metrics along the way. Classify errors: poison messages get quarantined and acked immediately, and only transient errors get nacked.

Treating every failure as poison. The opposite mistake is worse. Acking a message after a database timeout because "we logged it" is silent data loss. When in doubt, nack and let the dead-letter topic catch what never succeeds.

Dedupe TTL shorter than redelivery horizon. Covered above, but it's the bug I see most often in review: a 24-hour Redis TTL on a subscription with 7-day retention, and a replay after an incident that re-applies a week of effects.

Consumer-level dedupe keys shared across consumers. If two subscriptions both write processedEvents/{eventId}, the first consumer to process an event makes the second one skip it. Always namespace the key by consumer.

Pub/Sub exactly-once and push. If you move a consumer from pull to push to save on idle instances, you lose exactly-once delivery, because it isn't available on push. Your idempotency guard needs to handle that on its own, which is one more reason to never rely on the broker feature by itself.

Checklist#

Checklist

  • Every event carries a stable ID generated before the publish call, outside any retry loop
  • The idempotency key is a natural business key where the domain has one, otherwise the event ID
  • Keys are namespaced per consumer
  • Dedupe record commits in the same transaction as the effect whenever the effect is in the same database
  • External effects use a claim with a lease plus the provider's idempotency key
  • Counters are guarded increments, sharded if hot, or derived from set membership
  • Push handler returns 2xx for success, duplicates and quarantined poison; non-2xx only for transient failures
  • Subscription has a retry policy with exponential backoff
  • Dead-letter topic is configured, has its own subscription, and alerts on backlog
  • Pub/Sub service agent has publisher on the DLQ and subscriber on the source subscription
  • Dedupe TTL exceeds message retention plus any replay window
  • Ordered consumers carry a version or sequence number and reject stale writes

When not to do this#

If the effect is already naturally idempotent (overwriting a document with the latest profile snapshot, recomputing a cache entry from source, setting a status to "shipped"), a dedupe store adds cost and contention and gives you nothing. Model the write as "make it so" and skip the guard. The same applies to best-effort telemetry, where a duplicate event within some percentage points of noise doesn't change a decision. Save the machinery for effects that people or money can see: counters users look at, messages they receive, and anything that moves funds or calls a partner.

Share
All articles →