Job Lifecycle Events and Hooks
In this tutorial, you will learn about Job Lifecycle Events and Hooks. We cover key concepts, practical examples, and best practices to help you master this topic.
Track job lifecycle events from creation to completion with event hooks for enqueue, start, success, failure, retry, and completion stages in background processing.
What You Learn
You will learn how to implement lifecycle events, register event listeners, build event-driven workflows, and use events for auditing, metrics, and notifications.
Why It Matters
Lifecycle events decouple job processing from side effects. Instead of adding audit logging, metrics, and notifications to every job, you emit events and let listeners handle them.
Real-World Use
DodaTech's job system emits lifecycle events to: record audit trail in database, update real-time dashboard via Websocket, send Slack notification on failure, and collect Prometheus metrics. Each listener is independent.
Lifecycle Event Model
flowchart LR
C[Created] --> E[Enqueued]
E --> S[Started]
S --> P[Processing]
P --> SC[Success]
P --> F[Failure]
P --> R[Retry]
SC --> D[Done]
F --> D
R --> E
D --> A[Archived]
Event Emitter
import time
import threading
class EventEmitter:
def __init__(self):
self._listeners = {}
def on(self, event, handler):
if event not in self._listeners:
self._listeners[event] = []
self._listeners[event].append(handler)
return handler
def off(self, event, handler):
if event in self._listeners:
self._listeners[event] = [
h for h in self._listeners[event] if h != handler
]
def emit(self, event, **data):
listeners = self._listeners.get(event, [])
for handler in listeners:
try:
handler(data)
except Exception as e:
print(f"Listener error for {event}: {e}")
class LifecycleJob(EventEmitter):
def __init__(self, name):
super().__init__()
self.name = name
def enqueue(self, job_data):
self.emit('enqueued', job_name=self.name, data=job_data, timestamp=time.time())
def start(self):
self.emit('started', job_name=self.name, timestamp=time.time())
def success(self, result):
self.emit('succeeded', job_name=self.name, result=result, timestamp=time.time())
def failure(self, error):
self.emit('failed', job_name=self.name, error=str(error), timestamp=time.time())
def retry(self, attempt):
self.emit('retried', job_name=self.name, attempt=attempt, timestamp=time.time())
job = LifecycleJob('email_sender')
@job.on('enqueued')
def log_enqueue(data):
print(f"Log: Job enqueued - {data['job_name']}")
@job.on('started')
def metrics_start(data):
print(f"Metrics: Job started - {data['job_name']}")
@job.on('succeeded')
def notify_success(data):
print(f"Notify: Job completed - {data['job_name']}")
@job.on('failed')
def log_failure(data):
print(f"Alert: Job FAILED - {data['job_name']}: {data['error']}")
job.enqueue({'to': 'user@example.com'})
job.start()
job.success({'sent': True})
Expected output:
Log: Job enqueued - email_sender
Metrics: Job started - email_sender
Notify: Job completed - email_sender
Lifecycle-Managed Worker
import time
import threading
import random
class LifecycleWorker:
def __init__(self):
self.listeners = {
'before_enqueue': [],
'after_enqueue': [],
'before_process': [],
'after_success': [],
'after_failure': [],
'after_retry': [],
}
def on(self, event, handler):
if event in self.listeners:
self.listeners[event].append(handler)
def _trigger(self, event, context):
for handler in self.listeners.get(event, []):
try:
handler(context)
except Exception as e:
print(f"Event handler error ({event}): {e}")
def process(self, job_func, **context):
self._trigger('before_enqueue', context)
print(f"Enqueued: {context.get('name', 'job')}")
self._trigger('after_enqueue', context)
self._trigger('before_process', context)
max_retries = context.get('max_retries', 3)
for attempt in range(max_retries):
try:
print(f" Attempt {attempt + 1}")
result = job_func(context)
context['result'] = result
self._trigger('after_success', context)
return result
except Exception as e:
context['error'] = str(e)
context['attempt'] = attempt + 1
if attempt < max_retries - 1:
self._trigger('after_retry', context)
time.sleep(0.5)
else:
self._trigger('after_failure', context)
raise
worker = LifecycleWorker()
@worker.on('before_enqueue')
def audit_enqueue(ctx):
print(f" Audit: Enqueuing {ctx['name']}")
@worker.on('after_success')
def metrics_success(ctx):
print(f" Metrics: {ctx['name']} succeeded")
@worker.on('after_failure')
def alert_failure(ctx):
print(f" Alert: {ctx['name']} failed after {ctx.get('attempt', 0)} attempts")
@worker.on('after_retry')
def log_retry(ctx):
print(f" Log: Retrying {ctx['name']} (attempt {ctx['attempt']})")
def unreliable_task(ctx):
if random.random() < 0.7:
raise ValueError("Transient error")
return {'processed': True}
try:
worker.process(unreliable_task, name='data_sync', max_retries=3)
except ValueError:
print("Job failed permanently")
Expected output:
Audit: Enqueuing data_sync
Enqueued: data_sync
Attempt 1
Log: Retrying data_sync (attempt 1)
Attempt 2
Metrics: data_sync succeeded
Event Persistence and Replay
import json
import time
import sqlite3
class EventStore:
def __init__(self, db_path=':memory:'):
self.conn = sqlite3.connect(db_path)
self.conn.execute('''
CREATE TABLE IF NOT EXISTS lifecycle_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
job_id TEXT,
event_type TEXT,
data TEXT,
created_at REAL
)
''')
self.conn.commit()
self.listeners = {}
def on(self, event_type, handler):
if event_type not in self.listeners:
self.listeners[event_type] = []
self.listeners[event_type].append(handler)
def record(self, job_id, event_type, data=None):
self.conn.execute(
'INSERT INTO lifecycle_events (job_id, event_type, data, created_at) VALUES (?, ?, ?, ?)',
(job_id, event_type, json.dumps(data or {}), time.time())
)
self.conn.commit()
for handler in self.listeners.get(event_type, []):
try:
handler(job_id, data)
except Exception as e:
print(f"Handler error: {e}")
def get_events(self, job_id=None, event_type=None):
query = 'SELECT * FROM lifecycle_events WHERE 1=1'
params = []
if job_id:
query += ' AND job_id = ?'
params.append(job_id)
if event_type:
query += ' AND event_type = ?'
params.append(event_type)
query += ' ORDER BY created_at'
cursor = self.conn.execute(query, params)
return [
{
'id': row[0], 'job_id': row[1],
'event_type': row[2],
'data': json.loads(row[3]),
'created_at': row[4],
}
for row in cursor.fetchall()
]
store = EventStore()
store.on('failed', lambda jid, d: print(f"Alert: Job {jid} failed"))
store.on('completed', lambda jid, d: print(f"Log: Job {jid} completed"))
store.record('scan-001', 'started', {'file': 'doc.pdf'})
store.record('scan-001', 'completed', {'threats': 0})
store.record('scan-002', 'started', {'file': 'virus.exe'})
store.record('scan-002', 'failed', {'error': 'Timeout'})
events = store.get_events(job_id='scan-002')
print(f"Events for scan-002: {len(events)}")
for ev in events:
print(f" {ev['event_type']}: {ev['data']}")
Expected output:
Log: Job scan-001 completed
Alert: Job scan-002 failed
Events for scan-002: 2
started: {'file': 'virus.exe'}
failed: {'error': 'Timeout'}
Common Mistakes
1. Blocking Event Handlers
Slow event handlers block job processing. Run event handlers asynchronously or in separate threads.
2. Event Handler Exceptions
One failing handler should not break other handlers or the job itself. Wrap each handler in try/except.
3. Too Many Events
Firing an event for every minor state change creates noise. Define meaningful lifecycle stages: enqueued, started, succeeded, failed, retried.
4. No Event Ordering Guarantee
Asynchronous event handlers may Process events out of order. Include timestamps and sequence numbers for ordering.
5. Memory Leaks from Listeners
Listeners registered but never removed accumulate over time. Implement listener cleanup for short-lived jobs.
Practice Questions
1. What are lifecycle events in job processing?
Signals emitted at key stages of a job's life: enqueued, started, succeeded, failed, retried. Listeners react to these events.
2. How do events differ from middleware?
Middleware wraps execution and can modify behavior. Events are notifications after the fact. Events are for side effects, middleware for cross-cutting concerns.
3. Why persist lifecycle events?
For auditing, debugging, replay, and analytics. Persisted events provide a complete history of every job's journey through the system.
4. How do you prevent event handler failures from affecting jobs?
Wrap each handler in try/except. Log handler failures but do not propagate them to the job execution.
Challenge
Build a lifecycle event system for a payment processing queue: events for created, validated, processing, succeeded, failed, refunded. Persist events to database. Event handlers for audit log, metrics, Slack notification, and customer email.
FAQ
Mini Project: Lifecycle System
import time
import json
class JobLifecycle:
def __init__(self):
self.handlers = {}
def on(self, event):
def decorator(func):
if event not in self.handlers:
self.handlers[event] = []
self.handlers[event].append(func)
return func
return decorator
def emit(self, event, **data):
for handler in self.handlers.get(event, []):
handler(data)
def run(self, job_func, **context):
job_id = context.get('id', str(time.time()))
self.emit('start', job_id=job_id, **context)
try:
result = job_func(context)
self.emit('success', job_id=job_id, result=result)
return result
except Exception as e:
self.emit('fail', job_id=job_id, error=str(e))
raise
lc = JobLifecycle()
@lc.on('start')
def on_start(data):
print(f" [Event] Started: {data['job_id']}")
@lc.on('success')
def on_success(data):
print(f" [Event] Succeeded: {data['job_id']}")
@lc.on('fail')
def on_fail(data):
print(f" [Event] Failed: {data['job_id']}: {data['error']}")
lc.run(lambda ctx: time.sleep(0.1) or 'ok', id='job-001')
Expected output:
[Event] Started: job-001
[Event] Succeeded: job-001
What's Next
Now that you understand lifecycle events, explore job metrics with Prometheus for monitoring, then learn about job dashboard for visualization.
Built by the developers of DodaTech
Doda Browser, DodaZIP & Durga Antivirus Pro