Structured Logging for Workers
Logs from background workers are the first thing an engineer reaches for when a job misbehaves and the last thing anyone designs, and this guide treats them as a first-class signal within Observability & Monitoring for Job Queues. Metrics tell you that the failure rate of send_invoice doubled; traces tell you where one slow job spent its time; logs tell you which invoice failed, on which attempt, with what error, and what the job did before it failed. That only works if every log line carries the same machine-readable context, and if the volume stays affordable when a busy fleet emits millions of lines an hour.
Request logs have a natural correlation key — the request id — and frameworks add it automatically. Job logs do not. A job may run minutes after the request that enqueued it, on a different host, several times, and its log lines are interleaved with those of dozens of other jobs on the same worker. Without deliberate structure, "show me everything that happened to invoice 8812" is a grep across hosts and hours.
The Scenario: Five Workers, One Missing Email
A support ticket reports that a customer never received their order confirmation. The engineer on call knows the order id and starts searching. The web logs show the order was created and a send_confirmation job enqueued. The worker logs are plain text: INFO Sending email to customer, WARNING SMTP timeout, retrying, INFO Email sent — thousands of each per hour, with no order id, no job id, and no indication of which lines belong together. Forty minutes later, the engineer finds the cause by correlating timestamps by hand: the job exhausted its retries during a mail-provider incident, and its final failure was logged at WARNING like every transient retry, so no alert fired.
With structured logs, the same investigation is one query: job_name="send_confirmation" AND order_id="8812", returning five lines — enqueued, attempt 1 failed, attempts 2–4 failed, attempt 5 failed permanently at ERROR — each with the job id, attempt number, queue, worker host, error class, and the trace id linking back to the request.
Architectural Overview: The Worker Log Pipeline
A worker log pipeline has four stages, and each has one job:
- Context binding in the worker: when a job starts, bind its identity (job id, name, queue, attempt) and relevant business ids to a logging context, so every line emitted during the job carries them automatically.
- Formatting: render each record as one JSON object per line on stdout. No multi-line tracebacks; the exception goes into a field.
- Collection: the platform's agent (Promtail, Grafana Alloy, Fluent Bit, Vector, the Datadog agent) tails container stdout and adds infrastructure labels — namespace, pod, node.
- Storage and query: Loki, Elasticsearch/OpenSearch, or a vendor, indexed on a small set of low-cardinality labels, with everything else searchable in the body.
{"ts":"2026-09-18T12:04:31.882Z","level":"error","msg":"job failed permanently",
"service":"notifications","job_name":"send_confirmation","job_id":"4f1c9e2a-7b3d",
"queue":"email","attempt":5,"max_attempts":5,"order_id":"8812","customer_id":"c-311",
"duration_ms":30012,"error_class":"SMTPTimeout","error":"timed out after 30s",
"trace_id":"7a1f3c9e4b2d8f60a1c3e5f7b9d1e3a5","span_id":"c3e5f7b9d1e3a5f7",
"worker":"notifications-worker-6b9f-x2k8p","enqueued_at":"2026-09-18T11:58:02.114Z"}
Two fields in that line deserve attention. attempt and max_attempts make "was this the last try?" answerable without a join. enqueued_at lets you compute queue wait from logs alone — useful when metrics are aggregated too coarsely to explain one job. The broader signal design is in Observability & Monitoring for Job Queues.
Implementation 1: Binding Job Context Automatically
The mechanism differs by language, but the pattern is the same: a hook that runs before every job binds context, and one that runs after clears it. Handlers then log normally.
# Celery + structlog: bind job context in signals, clear it afterwards
import structlog
from celery.signals import task_prerun, task_postrun, task_retry, task_failure
structlog.configure(
processors=[
structlog.contextvars.merge_contextvars, # pulls in bound job context
structlog.processors.add_log_level,
structlog.processors.TimeStamper(fmt="iso", utc=True),
structlog.processors.dict_tracebacks, # exception -> structured field
structlog.processors.JSONRenderer(),
],
)
log = structlog.get_logger()
@task_prerun.connect
def bind_job(task_id, task, args, kwargs, **_):
req = task.request
structlog.contextvars.clear_contextvars()
structlog.contextvars.bind_contextvars(
job_id=task_id, job_name=task.name, queue=(req.delivery_info or {}).get("routing_key"),
attempt=req.retries + 1, max_attempts=(task.max_retries or 0) + 1,
**{k: v for k, v in kwargs.items() if k in ("order_id", "customer_id", "tenant_id")},
)
@task_postrun.connect
def clear_job(**_):
structlog.contextvars.clear_contextvars() # thread reuse: never leak context
Only whitelisted business ids are bound — never whole argument dicts, which leak personal data into logs and bloat every line. The Celery-specific details, including propagating the enqueuing request's correlation id, are in correlating logs with job IDs in Celery; the Node equivalent with pino child loggers is in structured JSON logging for BullMQ workers.
Implementation 2: Log Levels That Mean Something for Retries
Retries are where most job logging goes wrong. A transient failure that will be retried is not an error — logging it at ERROR trains everyone to ignore errors. A failure on the final attempt is an error, and logging it at WARNING hides exactly the event that matters. Make the level a function of the attempt:
def log_job_failure(exc: Exception, attempt: int, max_attempts: int, will_retry: bool) -> None:
fields = {"error_class": type(exc).__name__, "error": str(exc)[:500]}
if not will_retry:
log.error("job failed permanently", **fields, final=True) # alertable
elif attempt >= max_attempts - 1:
log.warning("job failed, last retry scheduled", **fields)
else:
log.info("job failed, will retry", **fields) # expected noise
@task_failure.connect
def on_failure(task_id, exception, **_):
req = current_task.request
will_retry = False # task_failure fires after retries end
log_job_failure(exception, req.retries + 1, current_task.max_retries + 1, will_retry)
@task_retry.connect
def on_retry(request, reason, **_):
log_job_failure(reason, request.retries + 1, current_task.max_retries + 1, True)
With this convention, a log-based alert on level="error" AND final=true fires once per job that actually gave up, and never for the thousands of retries that later succeeded. It pairs with the metrics-based alerts in alerting on stuck and stalled jobs.
Trade-off Analysis: What Goes in Labels, Fields, and Nowhere
| Field | Loki label / ES keyword | Body field | Never log |
|---|---|---|---|
| service, environment | Label | — | — |
| queue, job_name | Label (bounded set) | — | — |
| level | Label | — | — |
| job_id, trace_id | — | Field (high cardinality) | — |
| order_id, customer_id, tenant_id | — | Field | — |
| attempt, duration_ms, error_class | — | Field | — |
| Full job arguments | — | — | Yes: PII, secrets, size |
| Email addresses, tokens, card data | — | — | Yes |
In Loki, labels define streams and every distinct label combination is a stream; putting job_id in a label creates one stream per job and degrades the whole cluster. In Elasticsearch, high-cardinality fields are fine as keywords, but mapping explosion from arbitrary argument keys is not — whitelist fields rather than dumping dicts. The Loki-specific setup is in shipping worker logs to Loki.
Failure Modes & Recovery
Context leaking between jobs. Worker threads and processes are reused. If context is bound at job start but not cleared at the end, the next job's lines carry the previous job's ids — worse than no ids, because they mislead. Always clear in a finally-equivalent hook, and test it: run two jobs in sequence on one thread and assert the second job's lines carry only its own ids.
Multi-line tracebacks. A Python or Java stack trace printed to stdout becomes twenty separate log lines, none with the job context. Render exceptions into a single field (dict_tracebacks, pino's err serializer, Logback's JSON encoder) so the traceback travels with the line that explains it.
Log volume explosions. A retry storm or a tight loop can multiply log volume a hundredfold in minutes, driving up cost and sometimes back-pressuring the collection agent. Rate-limit or sample repetitive lines at the source; log sampling for high-volume queues covers which lines are safe to sample and which must never be.
Sensitive data in arguments. Job payloads routinely contain emails, addresses, and tokens. Logging kwargs wholesale ships them to a system with broad access and long retention. Whitelist the ids you bind and add a redaction processor as a backstop.
Performance Tuning
- Log less per job, not less context. One line at start (at
debug, sampled), one at completion with duration and outcome, and failures. Handlers that log every loop iteration produce most of the volume and little of the value. - Use asynchronous or buffered handlers carefully. They cut latency in hot paths but lose the last lines on a crash — precisely the lines you need. Flush on shutdown signals, and prefer stdout with the platform agent handling buffering.
- Keep lines small. Truncate error messages (500 characters is plenty), never log payloads, and avoid nested objects that the backend must flatten.
- Serialize once. Structured loggers that render JSON per handler (one for stdout, one for a file, one for a vendor SDK) multiply CPU cost on hot workers. Emit once to stdout and let the collection agent fan out to destinations; a profiler on a busy worker often shows JSON rendering in the top ten functions when this is not done.
- Put derived metrics in metrics. Counting failures by parsing logs is slow and expensive. Emit counters and histograms from the same hooks that log, as in instrumenting Celery with a Prometheus exporter, and use logs for the details.
# Loki: jobs that failed permanently in the last hour, by job name
sum by (job_name) (count_over_time({service="notifications"} | json | level="error" | final="true" [1h]))
# Full history of one business entity across retries and workers
{service="notifications"} | json | order_id="8812"
Connecting Logs to Traces and Metrics
Structured logs are most useful when they link sideways to the other two signals. The link to traces is the trace_id field: if the enqueuing request's trace context is propagated into the job (as covered in propagating trace context through Celery tasks), the logging hook can read the active span and add its ids to every line. Grafana, Datadog, and most vendors then render a "view trace" link next to each log line, and a "view logs" link on each span.
from opentelemetry import trace
def add_trace_ids(logger, method, event_dict):
span = trace.get_current_span()
ctx = span.get_span_context()
if ctx.is_valid:
event_dict["trace_id"] = format(ctx.trace_id, "032x")
event_dict["span_id"] = format(ctx.span_id, "016x")
return event_dict
structlog.configure(processors=[structlog.contextvars.merge_contextvars, add_trace_ids,
structlog.processors.add_log_level,
structlog.processors.TimeStamper(fmt="iso", utc=True),
structlog.processors.JSONRenderer()])
The link to metrics is the label set: use the same job_name and queue values in log labels and metric labels. An alert on job_failures_total{job_name="send_confirmation"} can then link straight to a log query filtered on the same job name, and the on-call engineer moves from "the rate went up" to "these are the failing jobs" in one click. Consistent naming across signals is a small discipline with an outsized effect on incident time.
A Shared Schema Across Languages
Most organisations run workers in more than one language — Celery in the data team, BullMQ in the product backend, Sidekiq in the billing service, a Go consumer somewhere. If each emits its own field names (jobId, job_id, jid, task_id), every cross-service query needs a translation table and every dashboard is per-language. Agree on a small schema once and enforce it in each language's logging setup.
# log-schema.yaml — the contract every worker's logger must satisfy
required:
ts: RFC 3339 UTC timestamp with milliseconds
level: one of debug | info | warning | error
msg: short, constant message (no interpolated ids)
service: deployable name, e.g. "notifications"
job_name: framework task/job/kind name
job_id: framework job id
queue: queue name the job was taken from
attempt: 1-based attempt number
optional:
max_attempts, duration_ms, enqueued_at, final, error_class, error,
trace_id, span_id, tenant_id, order_id, customer_id
forbidden:
args, kwargs, payload, email, token, password, card
Two rules in that contract carry most of the value. msg is a constant string — "job failed permanently", never "job 4f1c for order 8812 failed" — so messages can be grouped and counted, while the ids live in fields. And the forbidden list is enforced mechanically: a redaction processor drops those keys even if a developer binds them by accident. A contract test in each service's suite renders one log line through the real logging configuration and validates it against the schema, which catches drift the day it is introduced rather than during the next incident.
Investigating an Incident with Worker Logs
Structured logs pay off in a predictable sequence of queries. Taking the missing-email scenario as an example, an investigation typically moves from broad to narrow in four steps, each a single query:
# 1. Is something failing more than usual? (volume of final failures by job)
sum by (job_name) (count_over_time({service="notifications"} | json | final="true" [30m]))
# 2. What is failing? (top error classes for the affected job)
topk(5, sum by (error_class) (count_over_time(
{service="notifications", job_name="send_confirmation"} | json | level="error" [30m])))
# 3. Who is affected? (distinct business ids among final failures)
{service="notifications", job_name="send_confirmation"} | json | final="true"
| line_format "{{.order_id}} {{.customer_id}} {{.error_class}}"
# 4. What happened to one of them? (full history across attempts and workers)
{service="notifications"} | json | order_id="8812"
Step three produces the list support needs to contact or remediate affected customers; step four, combined with the trace_id link, reaches the original request. None of it is possible without the fields bound in Implementation 1, and all of it is fast because the labels (service, job_name) narrow the streams before any JSON is parsed.
Retention and Access
Job logs contain business identifiers, and sometimes more than intended. Treat retention and access as part of the design rather than an afterthought. Keep full-detail worker logs for as long as incidents are realistically investigated — typically 14 to 30 days — and archive or drop them afterwards; long-term trends belong in metrics, which are far cheaper to keep. Restrict who can query logs from workers that process payments or personal data, and make sure deletion requests for personal data can reach log storage if your regulatory environment requires it. The whitelist-and-redact approach in Implementation 1 keeps the sensitive surface small enough that these policies are practical.
FAQ
Should job logs and request logs share a schema? Yes, for the common fields — timestamp, level, service, trace_id, message — so dashboards and queries work across both. Job logs add job_id, job_name, queue, and attempt.
How do I find the request that enqueued a job? Propagate the request's trace id or correlation id into the job's headers at enqueue time and bind it in the job's context. The trace id links logs and traces in one step.
Is it worth logging successful jobs? One completion line per job with duration is valuable at low to moderate volume. At very high volume, sample successes and keep every failure.
What about workers written without a structured logger?
Wrap the existing logger rather than rewriting call sites. In Python, a logging.Formatter subclass that emits JSON and reads bound context from contextvars makes every existing logger.info(...) call structured. In Ruby, Sidekiq's Sidekiq.logger.formatter = Sidekiq::Logger::Formatters::JSON.new switches its own lines to JSON, and Rails' tagged logging can carry job ids. Migrate the highest-volume and most-investigated jobs first.
Do I still need logs if I have traces?
Yes, for different reasons. Traces are usually sampled and show timing and structure; logs are usually complete for failures and carry the business detail — the exact error text, the ids involved, the decision a handler made. The two work best linked by trace_id rather than as alternatives.
Where should log-based alerts point? At a saved query that shows the failing jobs with their context, so the alert's first click answers "which jobs, which entities, which error".
Related
- Correlating Logs with Job IDs in Celery — structlog, signals, and request correlation.
- Structured JSON Logging for BullMQ Workers — pino child loggers for Node workers.
- Log Sampling for High-Volume Queues — keep costs bounded without losing failures.
- Shipping Worker Logs to Loki — labels, pipelines, and queries.
- Distributed Tracing for Async Jobs — the signal logs link to through trace_id.