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.

From grep to one query On the left, plain-text lines from five workers are interleaved with no job id or business id, so finding one job's history means correlating timestamps by hand. On the right, every JSON line carries job id, job name, attempt, queue, and order id, so a single filter on order id returns the five lines describing that job from enqueue to final failure. Finding one job's story plain text INFO Sending email to customer WARNING SMTP timeout, retrying INFO Sending email to customer WARNING SMTP timeout, retrying which lines are order 8812? structured filter: order_id = 8812 attempt 1-4: WARNING, transient attempt 5: ERROR, retries exhausted 5 lines, one query

Architectural Overview: The Worker Log Pipeline

A worker log pipeline has four stages, and each has one job:

  1. 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.
  2. Formatting: render each record as one JSON object per line on stdout. No multi-line tracebacks; the exception goes into a field.
  3. Collection: the platform's agent (Promtail, Grafana Alloy, Fluent Bit, Vector, the Datadog agent) tails container stdout and adds infrastructure labels — namespace, pod, node.
  4. 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.

Bind, format, collect, store When a job starts, the worker binds job id, name, queue, attempt and business ids to the logging context. Each record is rendered as one JSON line on stdout. A node agent such as Promtail or Fluent Bit tails container output and adds namespace and pod labels. Loki or Elasticsearch stores the lines, indexed on a few labels with the rest searchable. Worker log pipeline bind context job id, attempt, ids JSON to stdout one line per record agent adds pod, namespace Loki / ES few labels, full text Context is bound once per job; every line the handler emits inherits it without extra code.

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.

Level follows consequence A job with five attempts fails four times and then permanently. Attempts one to three are logged at info because a retry is scheduled and most such jobs later succeed. Attempt four is logged at warning because only one retry remains. Attempt five, the permanent failure, is logged at error with final set to true, which is what alerts key on. Five attempts, three levels 1: info 2: info 3: info 4: warning 5: error, final Most retried jobs succeed at attempt 2 or 3; their failures stay at info and never page anyone. Alert on level=error AND final=true: one page per job that truly gave up.

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.

Clear context or it leaks into the next job Worker thread 3 runs job A for order 8812 and binds order_id 8812. Job A finishes but the context is not cleared. The same thread then runs job B for order 9001, which binds only job_id; its log lines still show order_id 8812, sending investigators to the wrong order. One thread, two jobs, stale context job A: order_id=8812 bound job B logs still say order_id=8812 fix: clear_contextvars() in the post-run hook, and a test that runs two jobs on one thread Misleading context is worse than missing context: it sends the investigation to the wrong record.

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