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¶
- Redis 6.2+ for
XAUTOCLAIM; Postgres for the transactional variant. - Consumer groups and redelivery, from Message Brokers & Event Streams.
- The topic overview, Distributed Locks & Coordination.
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.
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.
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.
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.
Related¶
- Distributed Locks & Coordination — up to the topic overview.
- Running singleton jobs across replicas — the same claim idea for scheduled work.
- Concurrent Execution & Worker Patterns — the section overview.