Job Deduplication and Idempotency — Complete Guide
In this tutorial, you will learn about Job Deduplication and Idempotency. We cover key concepts, practical examples, and best practices to help you master this topic.
Prevent duplicate job execution using idempotency keys, deduplication sets, and at-least-once vs exactly-once processing strategies for reliable queues.
What You Learn
You will learn how to implement job deduplication with idempotency keys, build deduplication sets in Redis, use database constraints for uniqueness, and choose between at-least-once and exactly-once delivery.
Why It Matters
Network failures, client retries, and worker crashes cause duplicate job submissions. Without deduplication, the same email is sent twice, the same payment is processed twice, and the same file is scanned twice. Deduplication prevents these costly mistakes.
Real-World Use
DodaTech's Webhook receiver sees the same event delivered multiple times due to HTTP retries. Each webhook carries an X-Idempotency-Key header. The receiver checks if this key was already processed before executing the handler, ensuring safe retries.
Idempotency Key Pattern
import redis
import json
import time
import uuid
r = redis.Redis()
class IdempotentProcessor:
def __init__(self, expiry=86400):
self.expiry = expiry
self.processed_key = 'processed_keys'
def process_once(self, idempotency_key, handler, *args):
dedup_key = f'dedup:{idempotency_key}'
if r.setnx(dedup_key, '1'):
r.expire(dedup_key, self.expiry)
try:
result = handler(*args)
r.hset(self.processed_key, idempotency_key, 'completed')
return result
except Exception as e:
r.delete(dedup_key)
raise e
else:
print(f"Duplicate detected: {idempotency_key}")
return None
def was_processed(self, idempotency_key):
return r.exists(f'dedup:{idempotency_key}')
def send_email(recipient, subject):
print(f"Email sent to {recipient}: {subject}")
return {'status': 'sent', 'recipient': recipient}
processor = IdempotentProcessor()
key1 = str(uuid.uuid4())
processor.process_once(key1, send_email, 'alice@example.com', 'Welcome')
processor.process_once(key1, send_email, 'alice@example.com', 'Welcome')
Expected output:
Email sent to alice: Welcome
Duplicate detected: ...
Database-Backed Deduplication
import sqlite3
import time
import uuid
class DatabaseDeduplicator:
def __init__(self, db_path=':memory:'):
self.conn = sqlite3.connect(db_path)
self.conn.execute('''
CREATE TABLE IF NOT EXISTS processed_jobs (
idempotency_key TEXT PRIMARY KEY,
status TEXT,
processed_at TIMESTAMP,
result TEXT
)
''')
self.conn.commit()
def try_process(self, key, handler, *args):
cursor = self.conn.execute(
'SELECT status FROM processed_jobs WHERE idempotency_key = ?',
(key,)
)
row = cursor.fetchone()
if row:
print(f"Already processed: {key} (status: {row[0]})")
return None
try:
result = handler(*args)
self.conn.execute(
'''INSERT INTO processed_jobs (idempotency_key, status, processed_at, result)
VALUES (?, ?, ?, ?)''',
(key, 'completed', time.time(), str(result))
)
self.conn.commit()
return result
except Exception as e:
self.conn.execute(
'''INSERT INTO processed_jobs (idempotency_key, status, processed_at, result)
VALUES (?, ?, ?, ?)''',
(key, 'failed', time.time(), str(e))
)
self.conn.commit()
raise e
def get_status(self, key):
cursor = self.conn.execute(
'SELECT status, processed_at FROM processed_jobs WHERE idempotency_key = ?',
(key,)
)
row = cursor.fetchone()
if row:
return {'status': row[0], 'processed_at': row[1]}
return None
dedup = DatabaseDeduplicator()
def process_payment(amount, account):
print(f"Payment of ${amount} to {account}")
return {'txn_id': str(uuid.uuid4())}
key = 'pmt-20260628-001'
dedup.try_process(key, process_payment, 100, 'acc-123')
dedup.try_process(key, process_payment, 100, 'acc-123')
Expected output:
Payment of $100 to acc-123
Already processed: pmt-20260628-001 (status: completed)
Redis Set Deduplication
import redis
import json
import time
r = redis.Redis()
class RedisSetDeduplicator:
def __init__(self, dedup_set='job_dedup', ttl=3600):
self.dedup_set = dedup_set
self.ttl = ttl
def is_unique(self, job_id, job_data):
member = f"{job_id}:{hash(json.dumps(job_data, sort_keys=True))}"
added = r.sadd(self.dedup_set, member)
if added:
r.expire(self.dedup_set, self.ttl)
return True
return False
def enqueue_if_unique(self, queue, job_id, job_data):
if self.is_unique(job_id, job_data):
r.lpush(queue, json.dumps(job_data))
return True
print(f"Deduplicated: {job_id}")
return False
def count_deduplicated(self):
return r.scard(self.dedup_set)
def clear(self):
r.delete(self.dedup_set)
dedup = RedisSetDeduplicator()
scan_job = {'file': 'document.pdf', 'scan_type': 'malware'}
dedup.enqueue_if_unique('scans', 'scan-doc.pdf', scan_job)
dedup.enqueue_if_unique('scans', 'scan-doc.pdf', scan_job)
altered_job = {'file': 'document.pdf', 'scan_type': 'full'}
dedup.enqueue_if_unique('scans', 'scan-doc.pdf', altered_job)
Expected output:
Deduplicated: scan-doc.pdf
Time-Window Deduplication
import redis
import time
import json
r = redis.Redis()
class TimeWindowDeduplicator:
def __init__(self, window_seconds=60):
self.window = window_seconds
def is_duplicate(self, job_type, resource_id):
key = f"dedup_window:{job_type}:{resource_id}"
exists = r.exists(key)
if not exists:
r.setex(key, self.window, '1')
return exists
def enqueue(self, queue, job_data):
job_type = job_data.get('type', 'default')
resource_id = job_data.get('resource_id', str(job_data))
if self.is_duplicate(job_type, resource_id):
print(f"Duplicate within window: {job_type}/{resource_id}")
return False
r.lpush(queue, json.dumps(job_data))
return True
dedup = TimeWindowDeduplicator(window_seconds=30)
dedup.enqueue('events', {'type': 'file_scan', 'resource_id': 'file-123'})
dedup.enqueue('events', {'type': 'file_scan', 'resource_id': 'file-123'})
time.sleep(31)
dedup.enqueue('events', {'type': 'file_scan', 'resource_id': 'file-123'})
Expected output:
Duplicate within window: file_scan/file-123
Idempotent Worker Design
import hashlib
import json
class IdempotentWorker:
def __init__(self, seen_store=None):
self.seen = seen_store or {}
def compute_job_hash(self, job):
fields = {k: v for k, v in job.items() if k != 'attempt'}
serialized = json.dumps(fields, sort_keys=True)
return hashlib.sha256(serialized.encode()).hexdigest()
def process(self, job, handler):
job_hash = self.compute_job_hash(job)
if job_hash in self.seen:
print(f"Skipping duplicate (hash: {job_hash[:8]}...)")
return self.seen[job_hash]
result = handler(job)
self.seen[job_hash] = result
return result
worker = IdempotentWorker()
def handle_webhook(payload):
print(f"Processing webhook: {payload['event']}")
return {'processed': payload['id']}
webhook = {'id': 'evt_123', 'event': 'user.created', 'data': {'name': 'Alice'}}
worker.process(webhook, handle_webhook)
worker.process(webhook, handle_webhook)
Expected output:
Processing webhook: user.created
Skipping duplicate (hash: a1b2c3d4...)
Common Mistakes
1. Relying on At-Most-Once Delivery
At-most-once can lose messages. Use at-least-once with idempotent processing. The consumer handles duplicates safely.
2. In-Memory Deduplication Only
In-memory sets are lost on restart. Use Redis or a database for persistent deduplication state.
3. No TTL on Dedup Keys
Without expiration, the deduplication set grows forever, consuming memory and slowing lookups. Always set a reasonable TTL.
4. Ignoring Job Content in Dedup Key
Deduplicating by job ID alone means different jobs with the same ID are incorrectly blocked. Include relevant fields in the dedup key.
5. Deduplication Without Logging
When a job is deduplicated, no record exists of why it was skipped. Log deduplication events for debugging and auditing.
Practice Questions
1. What is idempotency in job processing?
An operation is idempotent if executing it multiple times produces the same result as executing it once. The same job can be safely retried.
2. How does SETNX help with deduplication?
SETNX sets a key only if it does not exist. If it returns success, this is the first time. If it returns failure, a duplicate is detected.
3. What is the difference between at-least-once and exactly-once?
At-least-once guarantees delivery but may duplicate. Exactly-once prevents duplicates but is harder to achieve. Most systems use at-least-once with idempotent consumers.
4. Why add TTL to dedup keys?
To bound memory usage. A job that succeeded is unlikely to be retried after hours or days. A 24-hour TTL is usually sufficient.
Challenge
Build a deduplication system for a payment processing queue. Handle: same payment ID submitted multiple times (reject duplicates within 24 hours), same amount sent to same account within 60 seconds (rate limit), webhook retries with idempotency headers, and a read endpoint to check if a job was already processed.
FAQ
Mini Project: Deduplication System
import redis
import json
import time
import hashlib
r = redis.Redis()
class DedupSystem:
def __init__(self, ttl=86400):
self.ttl = ttl
def make_key(self, namespace, identifier):
raw = f"{namespace}:{identifier}"
return f"dedup:{hashlib.sha256(raw.encode()).hexdigest()}"
def check_and_mark(self, namespace, identifier):
key = self.make_key(namespace, identifier)
if r.setnx(key, '1'):
r.expire(key, self.ttl)
return True
return False
def enqueue_dedup(self, queue, job_data, namespace='job'):
identifier = job_data.get('id') or job_data.get('dedup_key')
if not identifier:
identifier = hashlib.md5(
json.dumps(job_data, sort_keys=True).encode()
).hexdigest()
if self.check_and_mark(namespace, identifier):
r.lpush(queue, json.dumps(job_data))
print(f"Enqueued: {identifier[:12]}...")
return True
else:
print(f"Duplicate blocked: {identifier[:12]}...")
return False
def enqueue_batch(self, queue, jobs):
enqueued = 0
blocked = 0
for job in jobs:
if self.enqueue_dedup(queue, job):
enqueued += 1
else:
blocked += 1
return enqueued, blocked
ds = DedupSystem()
ds.enqueue_dedup('email', {'id': 'email-welcome-1', 'to': 'a@x.com'})
ds.enqueue_dedup('email', {'id': 'email-welcome-1', 'to': 'a@x.com'})
ds.enqueue_dedup('email', {'id': 'email-report-1', 'to': 'b@x.com'})
Expected output:
Enqueued: email-welcom...
Duplicate blocked: email-welcom...
Enqueued: email-repor...
What's Next
Now that you understand job deduplication, explore job dependencies for chaining related tasks, then learn about distributed workers for scaling across multiple machines.
Built by the developers of DodaTech
Doda Browser, DodaZIP & Durga Antivirus Pro