Skip to content

Advanced Job Failure Handling — Complete Guide

DodaTech Updated 2026-06-28 7 min read

In this tutorial, you will learn about Advanced Job Failure Handling. We cover key concepts, practical examples, and best practices to help you master this topic.

Handle background job failures with circuit breakers, exponential backoff, jitter, dead letter queues, failure classification, and recovery workflows.

What You Learn

You will learn how to classify failures by type, implement circuit breakers for external service failures, build dead letter queues with replay capability, and design recovery workflows.

Why It Matters

Simple retry loops are not enough. Different failures need different responses: transient failures retry with backoff, permanent failures go to dead letter, systemic failures trigger circuit breakers.

Real-World Use

DodaTech's job processor classifies failures: network timeouts (retry 3x with backoff), invalid data (dead letter immediately), database unreachable (circuit breaker, retry later). Each type has distinct handling.

Failure Classification

flowchart TD
    J[Job Fails] --> C{Classify Error}
    C -->|Transient| T[Network/Timeout]
    C -->|Permanent| P[Invalid Data]
    C -->|Systemic| S[Service Down]
    T --> R{Retry Count}
    R -->|< Max| RB[Backoff Retry]
    R -->|= Max| DL[Dead Letter]
    P --> DL
    S --> CB[Circuit Breaker]
    CB -->|Open| W[Wait]
    W -->|Half-Open| RR[Re-test]
    RR -->|OK| CR[Close]
    RR -->|Fail| W

Failure Classifier

import time
import json

class FailureClassifier:
    TRANSIENT_ERRORS = [
        'ConnectionError', 'TimeoutError', 'ConnectionRefusedError',
        'socket.timeout', 'requests.exceptions.ConnectionError',
    ]
    PERMANENT_ERRORS = [
        'ValueError', 'TypeError', 'KeyError',
        'json.JSONDecodeError', 'ValidationError',
    ]
    SYSTEMIC_ERRORS = [
        'DatabaseUnavailable', 'ServiceUnavailable',
        'RedisConnectionError',
    ]

    @classmethod
    def classify(cls, error):
        error_name = type(error).__name__
        error_str = str(error)

        for pattern in cls.TRANSIENT_ERRORS:
            if pattern in error_name or pattern in error_str:
                return 'transient'
        for pattern in cls.PERMANENT_ERRORS:
            if pattern in error_name or pattern in error_str:
                return 'permanent'
        for pattern in cls.SYSTEMIC_ERRORS:
            if pattern in error_name or pattern in error_str:
                return 'systemic'

        return 'unknown'

    @classmethod
    def should_retry(cls, error, attempt, max_retries=3):
        failure_type = cls.classify(error)
        if failure_type == 'permanent':
            return False
        if failure_type == 'systemic':
            return attempt < 2
        if failure_type == 'transient':
            return attempt < max_retries
        return attempt < 2

def process_payment(amount):
    raise ConnectionError("Database connection refused")

try:
    process_payment(100)
except Exception as e:
    ftype = FailureClassifier.classify(e)
    retry = FailureClassifier.should_retry(e, 1)
    print(f"Failure type: {ftype}")
    print(f"Should retry: {retry}")

Expected output:

Failure type: transient
Should retry: True

Circuit Breaker for Job Failures

import time
import threading

class JobCircuitBreaker:
    def __init__(self, name, failure_threshold=5, recovery_timeout=30):
        self.name = name
        self.failure_threshold = failure_threshold
        self.recovery_timeout = recovery_timeout
        self.failure_count = 0
        self.state = 'closed'
        self.last_failure_time = 0
        self._lock = threading.Lock()

    def record_failure(self):
        with self._lock:
            self.failure_count += 1
            self.last_failure_time = time.time()
            if self.failure_count >= self.failure_threshold:
                self.state = 'open'
                print(f"[{self.name}] Circuit OPEN after {self.failure_count} failures")

    def record_success(self):
        with self._lock:
            if self.state == 'half-open':
                print(f"[{self.name}] Circuit CLOSED after successful test")
            self.failure_count = 0
            self.state = 'closed'

    def can_execute(self):
        with self._lock:
            if self.state == 'closed':
                return True
            if self.state == 'open':
                if time.time() - self.last_failure_time > self.recovery_timeout:
                    self.state = 'half-open'
                    print(f"[{self.name}] Circuit HALF-OPEN, testing")
                    return True
            return False

    def execute(self, func, *args, **kwargs):
        if not self.can_execute():
            raise Exception(f"Circuit breaker OPEN for {self.name}")

        try:
            result = func(*args, **kwargs)
            self.record_success()
            return result
        except Exception as e:
            self.record_failure()
            raise

def failing_api_call():
    raise TimeoutError("API timeout")

cb = JobCircuitBreaker('payment-api', failure_threshold=3, recovery_timeout=2)

for i in range(5):
    try:
        cb.execute(failing_api_call)
    except Exception as e:
        print(f"Attempt {i+1}: {e}")

time.sleep(3)
try:
    cb.execute(lambda: "API response")
except Exception as e:
    print(f"After recovery: {e}")

Expected output:

[payment-api] Circuit OPEN after 3 failures
Attempt 1: API timeout
Attempt 2: API timeout
Attempt 3: API timeout
Attempt 4: Circuit breaker OPEN for payment-api
Attempt 5: Circuit breaker OPEN for payment-api
[payment-api] Circuit HALF-OPEN, testing
[payment-api] Circuit CLOSED after successful test

Dead Letter Queue with Replay

import redis
import json
import time

r = redis.Redis()

class DeadLetterQueue:
    def __init__(self, dlq_name='dead_letter'):
        self.dlq_name = dlq_name

    def send_to_dlq(self, job, reason, original_queue=None):
        dlq_entry = {
            'original_job': job,
            'original_queue': original_queue,
            'failure_reason': reason,
            'failed_at': time.time(),
            'dlq_id': f"dlq-{time.time_ns()}",
        }
        r.lpush(self.dlq_name, json.dumps(dlq_entry))
        r.hincrby('dlq_stats', reason, 1)
        return dlq_entry['dlq_id']

    def replay_job(self, dlq_id, target_queue=None):
        all_entries = []
        while True:
            data = r.rpop(self.dlq_name)
            if not data:
                break
            all_entries.append(json.loads(data))

        replayed = None
        for entry in all_entries:
            if entry['dlq_id'] == dlq_id:
                queue = target_queue or entry.get('original_queue', 'default')
                r.lpush(queue, json.dumps(entry['original_job']))
                replayed = entry
                print(f"Replayed {dlq_id} to {queue}")
            else:
                r.lpush(self.dlq_name, json.dumps(entry))
        return replayed

    def replay_all(self, target_queue=None):
        count = 0
        while True:
            data = r.rpop(self.dlq_name)
            if not data:
                break
            entry = json.loads(data)
            queue = target_queue or entry.get('original_queue', 'default')
            r.lpush(queue, json.dumps(entry['original_job']))
            count += 1
        print(f"Replayed {count} jobs")
        return count

    def get_stats(self):
        stats = r.hgetall('dlq_stats')
        return {k.decode(): int(v) for k, v in stats.items()}

    def list_entries(self, limit=10):
        entries = []
        for i in range(limit):
            data = r.lindex(self.dlq_name, i)
            if data:
                entries.append(json.loads(data))
        return entries

dlq = DeadLetterQueue()
job = {'task': 'process_payment', 'amount': 100, 'account': 'acc-001'}
dlq_id = dlq.send_to_dlq(job, 'invalid_account', 'payments')
print(f"DLQ entry: {dlq_id}")
dlq.replay_all(target_queue='retry_queue')
stats = dlq.get_stats()
print(f"DLQ stats: {stats}")

Expected output:

DLQ entry: dlq-...
Replayed 1 jobs
DLQ stats: {'invalid_account': 1}

Recovery Workflow

import time
import json
import threading

class RecoveryWorkflow:
    def __init__(self):
        self.recovery_handlers = {}

    def on_failure(self, failure_type):
        def decorator(func):
            self.recovery_handlers[failure_type] = func
            return func
        return decorator

    def handle_failure(self, job, error):
        failure_type = FailureClassifier.classify(error)
        handler = self.recovery_handlers.get(failure_type)

        if handler:
            print(f"Running recovery for {failure_type}")
            return handler(job, error)
        else:
            print(f"No recovery handler for {failure_type}")
            return {'action': 'dead_letter'}

rw = RecoveryWorkflow()

@rw.on_failure('transient')
def handle_transient(job, error):
    retry_count = job.get('retry_count', 0) + 1
    if retry_count <= 3:
        backoff = 2 ** retry_count
        print(f"  Retry {retry_count} after {backoff}s backoff")
        return {'action': 'retry', 'delay': backoff, 'retry_count': retry_count}
    return {'action': 'dead_letter', 'reason': 'max_retries'}

@rw.on_failure('permanent')
def handle_permanent(job, error):
    print(f"  Permanent failure, sending to DLQ")
    return {'action': 'dead_letter', 'reason': str(error)}

@rw.on_failure('systemic')
def handle_systemic(job, error):
    print(f"  Systemic failure, circuit breaker")
    return {'action': 'circuit_break', 'timeout': 30}

job = {'task': 'email', 'retry_count': 0}
result = rw.handle_failure(job, TimeoutError("SMTP timeout"))
print(f"Recovery action: {result['action']}")

Expected output:

Running recovery for transient
  Retry 1 after 2s backoff
Recovery action: retry

Common Mistakes

1. Treating All Failures the Same

Network timeouts and invalid data need different handling. Classify failures and apply appropriate strategies.

2. Infinite Retries

Jobs that keep retrying forever waste resources. Always set max retry count and dead letter after exhaustion.

3. No Backoff Variation

Retrying at the same interval causes thundering herd when the service recovers. Use exponential backoff with jitter.

4. Ignoring Systemic Failures

If the database is down, retrying jobs will all fail. Use circuit breakers to stop retrying until the service recovers.

5. No Dead Letter Monitoring

Dead letter queues accumulate silently. Monitor DLQ depth and alert when jobs are being sent there.

Practice Questions

1. How do you classify job failures?

Transient (network, timeout, can retry), permanent (invalid data, cannot fix), systemic (service down, need circuit breaker).

2. What is a dead letter queue?

A queue for jobs that failed permanently. Jobs can be inspected, replayed, or discarded from the DLQ.

3. How does a circuit breaker help with failures?

It stops calling a failing service, giving it time to recover. Prevents cascading failures and wasted retries.

4. What is exponential backoff with jitter?

Retry delay increases exponentially (1s, 2s, 4s, 8s) with random variation (jitter) to prevent thundering herd.

Challenge

Build a failure handling system: classify failures as transient/permanent/systemic, retry transient with backoff and jitter, circuit break for systemic failures, dead letter for permanent failures, and alert when DLQ depth exceeds threshold.

FAQ

What is the difference between transient and permanent failures?

Transient failures resolve on retry (network glitch). Permanent failures will never succeed (invalid payload). Retry only transient.

How long should I wait before sending to dead letter?

After exhausting retries (3-5 attempts). The job has had fair chance to succeed. Further retries waste resources.

Can I recover jobs from dead letter?

Yes. DLQ allows replaying jobs to their original queue or a special retry queue. Inspect before replaying.

Should I use the same retry strategy for all jobs?

No. Payment jobs may retry more aggressively. Logging jobs may fail immediately. Configure per queue.

How do I prevent the thundering herd problem?

Use jitter in retry backoff. Randomize the delay so not all retries happen simultaneously when a service recovers.

Mini Project: Failure Handler

import time
import json
import random

class FailureHandler:
    def __init__(self, max_retries=3):
        self.max_retries = max_retries
        self.stats = {'transient': 0, 'permanent': 0, 'systemic': 0}

    def handle(self, job, error):
        error_name = type(error).__name__
        if error_name in ('ConnectionError', 'TimeoutError'):
            self.stats['transient'] += 1
            retry_count = job.get('_retries', 0)
            if retry_count < self.max_retries:
                delay = (2 ** retry_count) + random.uniform(0, 1)
                job['_retries'] = retry_count + 1
                return {'action': 'retry', 'delay': round(delay, 2)}
            return {'action': 'dead_letter'}
        elif error_name in ('ValueError', 'TypeError'):
            self.stats['permanent'] += 1
            return {'action': 'dead_letter'}
        else:
            self.stats['systemic'] += 1
            return {'action': 'circuit_break'}

fh = FailureHandler()
print(fh.handle({'task': 'a'}, ConnectionError("timeout")))
print(fh.handle({'task': 'a'}, ValueError("bad data")))
print(f"Stats: {fh.stats}")

Expected output:

{'action': 'retry', 'delay': 1.23}
{'action': 'dead_letter'}
Stats: {'transient': 1, 'permanent': 1, 'systemic': 0}

What's Next

Now that you understand failure handling, explore dead letter queues for failed job management, then learn about retry with backoff and jitter for optimal retry strategies.

Built by the developers of DodaTech

Doda Browser, DodaZIP & Durga Antivirus Pro