Skip to content

Job Deduplication and Idempotency — Complete Guide

DodaTech Updated 2026-06-28 6 min read

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

What is the difference between deduplication and idempotency?

Deduplication prevents the job from running twice. Idempotency ensures running it twice has the same effect. Deduplication is proactive; idempotency is defensive.

Can I deduplicate without external storage?

Not reliably. In-memory deduplication is lost on restart. You need Redis, a database, or a distributed cache for persistent dedup state.

How long should dedup keys live?

24-48 hours for most systems. Long enough to cover retry windows. Shorter TTL reduces memory usage but may let duplicates through for very late retries.

Does Celery support built-in deduplication?

No. Celery does not have built-in dedup. Implement it with a custom task base class that checks a dedup key before executing.

Can deduplication cause job loss?

Only if the dedup key is too broad. A failed job with the same key as a future valid job would be incorrectly deduplicated. Scope keys to job identity, not content.

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