Advanced Job Failure Handling — Complete Guide
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
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