Deduplicating Jobs with Redis SET NX Keys
A Redis key written with SET key value NX EX ttl is the most common deduplication primitive in job systems — cheap, fast, and shared by every worker. This guide shows how to use it correctly for both enqueue-side and execution-side dedup, as part of Exactly-Once vs At-Least-Once Delivery in Queue Fundamentals & Architecture. It also marks clearly where a TTL-based key stops being enough and a durable record is needed.
Problem Statement
A notification service enqueues a "push notification" job whenever a user receives a comment. Three sources of duplicates reach users: the web tier retries enqueue on timeouts, a comment edit re-fires the same event within seconds, and SQS occasionally delivers the same message twice. Users receive two or three identical notifications several times a day. You want at most one notification per comment per user, a dedup check that adds under a millisecond, protection against concurrent duplicate executions, and no permanent blocking if a worker crashes after claiming a job.
Prerequisites
- A Redis instance reachable from producers and workers (a dedicated or shared instance with
noeviction, or at least with dedup keys protected from eviction). - A deterministic identity for each unit of work — here
(user_id, comment_id, "push"). - A job framework or consumer where you can run code before the side effect (all of them).
- An understanding of how long duplicates can arrive after the original: seconds for retries, minutes for redeliveries.
Step 1 — Choose the Dedup Key and Its Window
The key must identify the effect, not the message. Two messages that should produce the same notification must produce the same key; two different notifications must not collide.
def dedup_key(user_id: int, comment_id: int, channel: str = "push") -> str:
# One notification per (user, comment, channel). Not per message id:
# an edit re-fires the event with a new message id but the same effect.
return f"dedup:notify:{channel}:{user_id}:{comment_id}"
The TTL is the dedup window: duplicates arriving after it expires will not be caught. Size it to cover the longest realistic delay between duplicates — producer retries (seconds), broker redeliveries (visibility timeout multiples), and replays from a dead-letter queue (hours or days, which a TTL key should not try to cover).
Step 2 — Deduplicate at Enqueue
The cheapest duplicate is one never enqueued. Before publishing, set the key with NX; publish only if the set succeeded.
import redis
r = redis.Redis.from_url(REDIS_URL, decode_responses=True)
def enqueue_notification(user_id: int, comment_id: int) -> bool:
key = dedup_key(user_id, comment_id)
if not r.set(key, "enqueued", nx=True, ex=86400):
return False # already enqueued in the last 24 h
try:
send_push.delay(user_id=user_id, comment_id=comment_id)
except Exception:
r.delete(key) # enqueue failed: allow a retry
raise
return True
The delete on failure matters: without it, a broker outage during enqueue would leave the key set and block the notification for 24 hours. Enqueue-side dedup handles producer retries and repeated events; it does not help against broker redelivery, which happens after enqueue.
Step 3 — Deduplicate at Execution with Claim and Complete
Execution-side dedup protects against redelivery and concurrent duplicates. A single "done" key set after the side effect is not enough: two workers can both check, both see nothing, and both send. Claim first, then mark completion.
CLAIM_TTL = 300 # longer than the job's worst-case runtime
DONE_TTL = 86400
def run_once(key: str, work) -> str:
# Try to claim. NX makes this atomic across all workers.
if not r.set(key, "running", nx=True, ex=CLAIM_TTL):
state = r.get(key)
return "duplicate-done" if state == "done" else "duplicate-running"
try:
work()
except Exception:
r.delete(key) # release so a retry can run
raise
r.set(key, "done", ex=DONE_TTL) # extend and mark complete
return "ran"
@app.task(bind=True, acks_late=True, max_retries=5)
def send_push(self, user_id: int, comment_id: int):
result = run_once(f"exec:{dedup_key(user_id, comment_id)}",
lambda: push_client.send(user_id, render(comment_id)))
if result == "duplicate-running":
raise self.retry(countdown=30) # the other copy may still fail
A duplicate that finds running retries later rather than returning success: the first copy might crash, and if the duplicate were acknowledged, nobody would send the notification. The claim TTL is the crash-recovery mechanism — if the claiming worker dies, the key expires and a retry can claim it.
Step 4 — Make Release Safe with a Token and Lua
The delete in Step 3 has a race: if worker A's work outlives CLAIM_TTL, the key expires, worker B claims it, and then A's failure path deletes B's claim. Store a unique token in the claim and release only if the token still matches, atomically.
import uuid
RELEASE = r.register_script("""
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('DEL', KEYS[1])
end
return 0
""")
COMPLETE = r.register_script("""
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('SET', KEYS[1], 'done', 'EX', ARGV[2])
end
return nil
""")
def run_once_safe(key: str, work) -> str:
token = f"running:{uuid.uuid4()}"
if not r.set(key, token, nx=True, ex=CLAIM_TTL):
return "duplicate-done" if r.get(key) == "done" else "duplicate-running"
try:
work()
except Exception:
RELEASE(keys=[key], args=[token])
raise
if COMPLETE(keys=[key], args=[token, DONE_TTL]) is None:
log.warning("claim expired before completion", key=key) # work may have run twice
return "ran"
The warning branch is honest about the limit: if the claim expired mid-work, another worker may have run the job concurrently. The fix is a claim TTL comfortably above the job's p99.9 duration, or extending the claim periodically for long jobs — the same lease idea as visibility timeouts.
Step 5 — Keep Dedup Keys from Being Evicted
A dedup key that Redis evicts under memory pressure silently stops deduplicating. On a Redis that also serves as a cache with allkeys-lru, dedup keys are just as likely to be evicted as cache entries.
# Check the policy on the instance holding dedup keys
redis-cli CONFIG GET maxmemory-policy
# Prefer a dedicated instance with noeviction, or volatile-* policies where
# only keys with TTLs are evicted AND cache keys also carry TTLs (then dedup keys still compete)
redis-cli INFO memory | grep -E 'used_memory_human|maxmemory_human'
redis-cli --scan --pattern 'dedup:*' | head -1000 | wc -l # rough key count sample
Estimate memory: each key is roughly 80–120 bytes including overhead. At 5 million notifications a day with a 24-hour TTL, that is about 500 MB of dedup keys — worth planning explicitly, as described in Redis maxmemory policy for queues.
Step 6 — Use a Durable Record Where TTLs Are Not Enough
For effects that must never repeat regardless of timing — payments, invoices, anything replayed from a DLQ days later — a Redis key with a TTL is the wrong tool. Record completion in the same database transaction as the effect, with a unique constraint, and keep the Redis key as a fast pre-check.
def charge_once(invoice_id: int) -> None:
fast_key = f"exec:charge:{invoice_id}"
if r.get(fast_key) == "done":
return # cheap path for hot duplicates
with db.transaction():
inserted = db.execute(
"INSERT INTO processed_effects (effect_key) VALUES (%s) ON CONFLICT DO NOTHING",
(f"charge:{invoice_id}",)).rowcount
if not inserted:
return # durable record says done
gateway.charge(invoice_id, idempotency_key=f"charge-{invoice_id}")
r.set(fast_key, "done", ex=86400)
The durable approach is covered fully in idempotent consumers with Postgres unique constraints.
Verification
def test_concurrent_duplicates_run_once(redis_client, push_fake):
import threading
barrier = threading.Barrier(8)
def worker():
barrier.wait()
run_once_safe("exec:dedup:test:1", lambda: push_fake.send(1))
threads = [threading.Thread(target=worker) for _ in range(8)]
[t.start() for t in threads]; [t.join() for t in threads]
assert push_fake.sent == 1
def test_crash_releases_after_ttl(redis_client, monkeypatch):
monkeypatch.setattr(mod, "CLAIM_TTL", 1)
redis_client.set("exec:k", "running:dead", nx=True, ex=1) # simulated crashed claimer
time.sleep(1.2)
assert run_once_safe("exec:k", lambda: None) == "ran"
In production, count outcomes (ran, duplicate-done, duplicate-running) as a metric. A steady trickle of duplicates is expected; a spike points at a producer retry storm or a visibility timeout that is too short.
Gotchas & Edge Cases
Redis failover loses recent keys. With asynchronous replication, keys written just before a primary fails may be missing on the new primary. Dedup is best-effort across failovers; combine with idempotent effects for anything critical.
Keys that are too broad. A key per user rather than per effect suppresses legitimate notifications. Include every field that distinguishes one intended effect from another.
Redis Cluster. Keys in Lua scripts must hash to one slot; single-key scripts as above are fine.
Clock-free by design. TTLs are enforced by Redis's clock, not the workers', which avoids skew problems — another reason to prefer Redis TTLs over timestamps stored in values.
FAQ
Should I use SETNX or SET NX EX?
SET key value NX EX ttl. The older SETNX plus a separate EXPIRE is not atomic; a crash between them leaves a key that never expires.
Can the job framework do this for me?
Often, yes: BullMQ job ids and deduplication options, Sidekiq Enterprise unique jobs, and Asynq's TaskID/Unique implement variants of this pattern. See deduplicating jobs in BullMQ and Sidekiq unique jobs and deduplication.
How long should the TTL be? Longer than the longest delay between duplicates you want to catch, and short enough to keep memory bounded — commonly 1–24 hours for notifications, minutes for enqueue-retry dedup.
Related
- Exactly-Once vs At-Least-Once Delivery — why duplicates are inevitable.
- Preventing Duplicate Job Execution with Idempotency — the broader idempotency toolkit.
- Idempotent Consumers with Postgres Unique Constraints — the durable counterpart to TTL keys.
- Kafka Exactly-Once Semantics with Transactions — transport-level exactly-once in Kafka.