Skip to content

Deduplicating Work Across Replicas

Queues deliver at least once. Redelivery after a consumer crash, producer retries and visibility timeouts all mean that the same job reaches your replicas more than once, and with several asyncio consumers processing concurrently, the copies are often handled at the same moment. Measured with a Redis Streams consumer group, three consumer processes of 20 concurrent tasks each and 3,000 jobs each delivered twice, twenty messages apart: with no deduplication, all 3,000 jobs ran their side effect twice. Checking a "done" marker before running and setting it afterwards still let 2,101 duplicates through, because both copies checked before either finished. Claiming the job with SET NX before running it removed every duplicate — but when one consumer was killed mid-run, 15 and 20 jobs were never done in two runs, because their claims outlived the crashed process. A leased claim with a separate done marker gave 0 duplicates and 0 missing jobs in the same crash test, and narrowed the remaining risk to a crash between the effect and the marker, which produced 13–20 duplicates when that gap was deliberately widened to 25 ms. Recording the job ID and its effect in one Postgres transaction gave 0 and 0 in every run. This guide implements and chooses between these.

Prerequisites

1. Measure your duplicate rate before choosing a fix

Make duplicates visible: record every execution of the side effect with the job ID, then count jobs that ran more than once and jobs that never ran. A test that injects duplicates and a crash gives a baseline in seconds:

SELECT count(DISTINCT job_id)                AS jobs_done,
       count(*) - count(DISTINCT job_id)      AS duplicates,
       3000 - count(DISTINCT job_id)          AS missing
FROM effects;

The test harness published every job twice, twenty messages apart, into a Redis stream, ran three consumers with XREADGROUP and had survivors reclaim messages idle for 2 s with XAUTOCLAIM. Measured without deduplication: 3,000 jobs done, 3,000 duplicates, none missing. That is the at-least-once contract working as designed — every job ran, some more than once — and it is the starting point every strategy below has to improve on without introducing jobs that never run.

Verify: you can count duplicates and missing jobs for a run of your own consumer with injected duplicates.

3,000 jobs, each delivered twice, three consumer replicas A grid of 5 rows by 4 columns. 3,000 jobs, each delivered twice, three consumer replicas strategy duplicates missing after a consumer crash guarantee none 3,000 0 at-least-once check done marker, then act 2,101 - racy claim with SET NX, then act 0 15-20 at-most-once leased claim + done marker 0 0 at-least-once, rare duplicates dedup key + effect in one transaction 0 0 exactly-once effect Redis Streams consumer group, 60 concurrent handlers; crash = SIGKILL of one consumer.

2. Do not check and then act

The intuitive fix — skip the job if it is marked done, mark it done afterwards — is a race between two concurrent handlers:

async def handle(job_id: int) -> None:
    if await r.exists(f"done:{job_id}"):        # both copies pass this check...
        return
    await send_receipt(job_id)                   # ...and both send
    await r.set(f"done:{job_id}", 1, ex=86400)

Measured: 2,101 of 3,000 jobs still ran twice. Two copies twenty messages apart were nearly always in flight at the same time across 60 concurrent handlers, and both saw no marker. The check only helps when copies arrive far apart — a redelivery an hour later — which is exactly the case that matters least. Any deduplication that works under concurrency must claim the job atomically before doing the work, so that exactly one handler wins.

Verify: your deduplication step is a single atomic operation — SET NX, INSERT ... ON CONFLICT, a compare-and-set — not a read followed by a write.

3. Claim atomically, but lease the claim

SET NX makes the claim atomic, and it removed every duplicate. But a claim set before the work and never expired is a promise the crashed consumer can no longer keep:

if await r.set(f"claim:{job_id}", 1, nx=True, ex=86400):    # one winner
    await send_receipt(job_id)
await r.xack("jobs", "g", message_id)

Measured with one of three consumers killed 1.5 s into the run: 15 jobs in one run and 20 in another were never done. The killed process had claimed them; the redelivered copies found the claim and were acknowledged as handled. That is at-most-once — safe for a notification you would rather skip than repeat, wrong for anything that must happen. The fix is a short claim that expires and a separate, durable done marker written after the work:

async def handle(job_id: int, message_id: bytes, me: str) -> None:
    if await r.exists(f"done:{job_id}"):
        await r.xack("jobs", "g", message_id)
        return
    if not await r.set(f"claim:{job_id}", me, nx=True, px=3000):
        return                                  # someone is working on it: do not ack
    await send_receipt(job_id)
    await r.set(f"done:{job_id}", 1, ex=86400)
    await r.xack("jobs", "g", message_id)

A copy that loses the claim leaves its message unacknowledged, so if the winner crashes, the message is reclaimed later and the expired claim lets someone retry. Measured in the same crash test, twice: 0 duplicates and 0 missing jobs.

Verify: after kill -9 of a consumer mid-run, every job is eventually done and the pending list drains to zero.

Two copies of one job under a leased claim A sequence of 5 messages between 3 participants. Two copies of one job under a leased claim consumer A Redis consumer B SET claim NX PX 3000 -> OK SET claim NX -> nil: no ack run the effect SET done; XACK reclaimed: EXISTS done -> ack If A crashes before SET done, the claim expires and B's copy runs the job.

4. Close the last gap with one transaction

The leased claim still has a window: if a consumer crashes after the effect but before writing the done marker, the claim expires and the job runs again. Measured by deliberately widening that gap to 25 ms: 13 and 20 duplicates across two crash runs. When the effect is a database write, remove the window by recording the job ID in the same transaction as the effect:

async def handle(pool, job_id: int) -> None:
    async with pool.acquire() as conn, conn.transaction():
        fresh = await conn.fetchval(
            "INSERT INTO processed(job_id) VALUES ($1) ON CONFLICT DO NOTHING RETURNING job_id",
            job_id,
        )
        if fresh is None:
            return                               # already done, by someone, atomically
        await conn.execute("INSERT INTO receipts(job_id, ...) VALUES ($1, ...)", job_id)

Either both rows commit or neither does; a concurrent copy blocks on the primary key until the first transaction finishes, then sees the conflict. Measured: 0 duplicates and 0 missing jobs without a crash and in two runs with a crash. When the effect is external — an email, a payment, a call to another service — no transaction spans it, and the remaining tool is an idempotency key the external system honours: pass the job ID as the key so a repeated call is recognised there, as in idempotency keys for safe async retries.

Verify: the dedup key and the effect are written in one transaction, or the external call carries an idempotency key derived from the job ID.

How should this job be deduplicated? A decision on What is the side effect with 4 outcomes. How should this job be deduplicated? What is the side effect? a write to the same database dedup key in the same transaction 0 duplicates, 0 missing external, accepts idempotency keys leased claim + key = job id duplicates absorbed downstream external, no idempotency leased claim + done marker rare duplicate on crash better skipped than repeated permanent claim at-most-once: 15-20 lost Exactly-once is achievable only where one transaction covers the effect.

5. Expire dedup records deliberately

Deduplication state grows with every job, so it must expire — but only after duplicates can no longer arrive. Size the retention from the longest redelivery delay in your system, not from convenience:

# processed(job_id int PRIMARY KEY, processed_at timestamptz NOT NULL DEFAULT now())
DEDUP_RETENTION = timedelta(days=7)      # > max queue retention + max retry horizon

async def prune(pool) -> int:
    status = await pool.execute(
        "DELETE FROM processed WHERE processed_at < now() - $1::interval", DEDUP_RETENTION
    )
    return int(status.split()[-1])

In Redis, the ex=86400 on the done marker is the same decision: a duplicate arriving after 24 hours runs again. Producer retries arrive within seconds, redelivery after a crash within the visibility or claim timeout, but replays from a stream's retained history or a dead-letter queue can arrive days later — the retention must cover the longest of these. Where volume makes per-job records expensive, deduplicate per batch or per offset instead, as consumers of partitioned logs do by committing offsets with their output, discussed in Async Data Pipelines.

Verify: dedup retention is documented and is longer than the maximum replay or redelivery horizon.

Verification

Work is deduplicated across replicas when:

  • A test injects duplicates and a crash, and counts duplicates and missing jobs.
  • Every claim is atomic — no check-then-act.
  • Claims are leased, and messages are acknowledged only after the work or a done marker.
  • Database effects share a transaction with the dedup key, and external effects carry an idempotency key.

Diagnostic Hook: count, per consumer, how many messages were skipped because the job was already done or claimed. A sudden rise means duplicate delivery has increased — a producer retrying, a consumer group rebalancing or a visibility timeout shorter than processing time — and the deduplication layer is absorbing it; zero skips for weeks means it may not be wired in at all.

Pitfalls & edge cases

  • Check-then-act. Measured: 2,101 of 3,000 jobs duplicated under concurrency.
  • Permanent claims before the work. Measured: 15–20 jobs never done after a crash.
  • The effect-to-marker gap. Measured: 13–20 duplicates when it was widened.
  • Dedup records that expire too soon. Late replays run again.

Frequently Asked Questions

How do I stop multiple replicas processing the same job?

Claim the job atomically before running it — SET NX with a lease in Redis, or an INSERT with a unique key in the same transaction as the effect in Postgres. Checking a done marker first and setting it afterwards still let 2,101 of 3,000 duplicates through in testing.

Can I get exactly-once processing with Redis Streams?

Delivery stays at-least-once. The effect can be exactly-once when it is a database write committed in the same transaction as a unique job ID, which gave 0 duplicates and 0 missing jobs with a consumer crash in testing.

Why did jobs go missing after adding deduplication?

A claim set before the work and never expired survives a crashed consumer, so redelivered copies are skipped. Lease the claim and write a separate done marker after the work.

How long should I keep deduplication keys?

Longer than the latest a duplicate can arrive: retry windows, redelivery timeouts and any replay from stream history or dead-letter queues.