Correlating Logs with Job IDs in Celery

Celery's default log format includes the task name and id only on the worker's own "received" and "succeeded" lines; everything your task code logs in between has no link to the task at all. This guide fixes that and carries a correlation id from the web request through every task it causes, as part of Structured Logging for Workers in Observability & Monitoring for Job Queues.

Problem Statement

A Django application logs every request with a request_id from middleware. Requests enqueue Celery tasks, which sometimes enqueue further tasks. When a user reports that an export never arrived, engineers can find the request by request_id but then lose the thread: the export task's logs, and the logs of the three chunk tasks it fanned out to, carry no shared identifier. With a prefork pool of 16 processes per worker and 20 workers, the relevant lines are scattered across 320 processes. You want every line emitted inside a task to carry the task id, task name, attempt, and the originating request_id, including tasks enqueued by other tasks.

Prerequisites

  • Celery 5.3+ and Python 3.10+ (for contextvars behaviour in the prefork and thread pools).
  • structlog 23+ configured for JSON output, or the standard logging module with a JSON formatter.
  • A request-id middleware in the web application that stores the id in a context variable.
  • Workers running with the prefork or threads pool (gevent/eventlet need the same approach with their own context handling).

Step 1 — Configure structlog to Merge Context Variables

structlog.contextvars stores bound values in a context variable, so each thread or greenlet has its own context and every logger call merges it automatically. Route the standard library's logging through the same processor chain so Celery's own lines and third-party libraries are structured too.

# logging_setup.py — imported by both Django settings and the Celery app
import logging, sys
import structlog

shared = [
    structlog.contextvars.merge_contextvars,
    structlog.processors.add_log_level,
    structlog.processors.TimeStamper(fmt="iso", utc=True),
    structlog.processors.dict_tracebacks,
]

structlog.configure(
    processors=shared + [structlog.stdlib.ProcessorFormatter.wrap_for_formatter],
    logger_factory=structlog.stdlib.LoggerFactory(),
    cache_logger_on_first_use=True,
)

handler = logging.StreamHandler(sys.stdout)
handler.setFormatter(structlog.stdlib.ProcessorFormatter(
    foreign_pre_chain=shared,                        # stdlib/Celery records get the context too
    processors=[structlog.stdlib.ProcessorFormatter.remove_processors_meta,
                structlog.processors.JSONRenderer()],
))
root = logging.getLogger()
root.handlers = [handler]
root.setLevel(logging.INFO)

Stop Celery from replacing this configuration when the worker starts by connecting to the setup_logging signal; otherwise Celery installs its own plain-text handler.

from celery.signals import setup_logging

@setup_logging.connect
def keep_our_logging(**_):
    import logging_setup  # noqa: F401 — already configured on import
One context, every logger When a task starts, its id, name, attempt and request id are bound into a context variable. Log calls from task code through structlog, Celery's own stdlib logging, and third-party libraries all pass through the same processor chain, which merges the bound context and renders one JSON line to stdout. All log sources share the bound context task code (structlog) Celery internals (stdlib) requests, boto3 (stdlib) merge_contextvars job_id, request_id, attempt one JSON line Without foreign_pre_chain, library and Celery lines would miss the task context entirely.

Step 2 — Put the Request Id into Task Headers at Enqueue

Task arguments are the wrong place for correlation data — they change the task's signature and end up in every retry and result. Celery message headers are designed for it. A before_task_publish signal handler copies the current request id into the headers of every task published from that context.

# correlation.py
import contextvars, uuid
from celery.signals import before_task_publish

request_id_var: contextvars.ContextVar[str | None] = contextvars.ContextVar("request_id", default=None)

@before_task_publish.connect
def add_correlation_header(headers=None, **_):
    rid = request_id_var.get()
    if rid and headers is not None:
        headers.setdefault("request_id", rid)    # parent task's id wins if already set

# Django middleware sets the variable for the duration of a request
class RequestIdMiddleware:
    def __init__(self, get_response):
        self.get_response = get_response
    def __call__(self, request):
        rid = request.headers.get("X-Request-ID") or uuid.uuid4().hex
        token = request_id_var.set(rid)
        structlog.contextvars.bind_contextvars(request_id=rid)
        try:
            return self.get_response(request)
        finally:
            request_id_var.reset(token)
            structlog.contextvars.clear_contextvars()

Because the handler runs at publish time in the publishing process, it works for tasks enqueued from web requests, management commands, and — importantly — other tasks, where the same context variable is set by Step 3.

Step 3 — Bind Task Context in task_prerun, Clear It in task_postrun

When the worker starts executing a task, read the header back, set the context variable (so child tasks inherit it through Step 2), and bind the task's identity for logging.

from celery.signals import task_prerun, task_postrun

@task_prerun.connect
def bind_task_context(task_id, task, **_):
    req = task.request
    rid = getattr(req, "request_id", None) or (req.headers or {}).get("request_id") or task_id
    request_id_var.set(rid)
    structlog.contextvars.clear_contextvars()
    structlog.contextvars.bind_contextvars(
        request_id=rid,
        job_id=task_id,
        job_name=task.name,
        attempt=req.retries + 1,
        queue=(req.delivery_info or {}).get("routing_key"),
        parent_job_id=req.parent_id,          # set when enqueued from another task
        root_job_id=req.root_id,              # first task in the chain/workflow
    )

@task_postrun.connect
def clear_task_context(**_):
    structlog.contextvars.clear_contextvars()
    request_id_var.set(None)

Celery already tracks parent_id and root_id for tasks started from other tasks and for canvas workflows. Binding them gives you a tree: filter by root_job_id to see a whole workflow — useful for the chains and chords in Celery chains, groups, and chords.

One request id across a task tree The web request with request id r-7c2 enqueues the export task; the before_task_publish handler copies r-7c2 into its headers. The worker binds r-7c2 when the export task starts. The export task enqueues three chunk tasks, and the same handler copies r-7c2 into their headers, while Celery sets their parent and root ids to the export task. request_id r-7c2 everywhere POST /exports request_id r-7c2 build_export header r-7c2 chunk 1: r-7c2, root=export chunk 2: r-7c2, root=export chunk 3: r-7c2, root=export Filter request_id="r-7c2" to see the request and all four tasks in one timeline.

Step 4 — Add Business Ids from Task Arguments

Job ids answer "what happened to this task"; business ids answer "what happened to this order". Bind a whitelisted set of argument names so support can search by the ids customers give them.

BUSINESS_KEYS = {"order_id", "customer_id", "export_id", "tenant_id"}

@task_prerun.connect
def bind_business_ids(task, args, kwargs, **_):
    ids = {k: v for k, v in (kwargs or {}).items() if k in BUSINESS_KEYS}
    # Positional args: map names from the task signature once, then pick whitelisted ones
    params = getattr(task, "_param_names", None)
    if params is None:
        params = list(inspect.signature(task.run).parameters)
        task._param_names = params
    ids.update({n: v for n, v in zip(params, args or ()) if n in BUSINESS_KEYS})
    structlog.contextvars.bind_contextvars(**ids)

Signal handlers run in registration order; register this after bind_task_context so the clear in the first handler does not wipe these values. Never bind whole argument dictionaries — they carry personal data and make every line large.

Step 5 — Log Outcomes with Duration and Final-Failure Flags

With context bound, outcome lines need only the outcome itself. Emit one completion line per task with duration, and failure lines whose level reflects whether retries remain.

import time
from celery.signals import task_success, task_retry, task_failure

log = structlog.get_logger("jobs")
_start: contextvars.ContextVar[float] = contextvars.ContextVar("job_start", default=0.0)

@task_prerun.connect
def start_timer(**_):
    _start.set(time.monotonic())

def _elapsed_ms() -> int:
    return int((time.monotonic() - _start.get()) * 1000)

@task_success.connect
def on_success(sender=None, **_):
    log.info("job succeeded", duration_ms=_elapsed_ms())

@task_retry.connect
def on_retry(request=None, reason=None, **_):
    log.info("job failed, will retry", duration_ms=_elapsed_ms(),
             error_class=type(reason).__name__, error=str(reason)[:500])

@task_failure.connect
def on_failure(exception=None, **_):
    log.error("job failed permanently", duration_ms=_elapsed_ms(), final=True,
              error_class=type(exception).__name__, error=str(exception)[:500])

The completion line is also the cheapest source of per-entity latency. Because it carries duration_ms next to order_id and job_name, a single query answers "how long did exports take for this tenant yesterday" without a dedicated metric, and the slowest individual jobs are one sort away. Keep it to one line per task: a start line, a completion line, and failure lines are enough; per-iteration logging inside the task is what turns a useful log stream into an expensive one.

Signals emit exactly one outcome line task_prerun binds the task and business context and starts a timer. When the task body finishes, exactly one of three signals fires: task_success logs job succeeded at info with the duration, task_retry logs job failed will retry at info, and task_failure logs job failed permanently at error with final set to true. task_postrun then clears the context. Signal lifecycle of one task task_prerun bind, start timer task_success: info, duration_ms task_retry: info, will retry task_failure: error, final=true task_postrun clear context

task_failure fires only when the task finally fails (retries exhausted or no retry configured), so level=error plus final=true marks jobs that truly gave up — the convention described in Structured Logging for Workers.

Verification

Run a request that fans out, then query by its id:

curl -s -H 'X-Request-ID: verify-123' -X POST https://staging.example.com/exports
# Loki: every line from the request and its four tasks
logcli query '{env="staging"} | json | request_id="verify-123"' --limit 200 \
  | jq -r '[.ts, .job_name // "web", .job_id // "-", .msg] | @tsv'

Expect web lines followed by build_export and three build_chunk tasks, all with request_id=verify-123, and chunk tasks carrying root_job_id equal to the export task's id. Add a unit test that runs two tasks sequentially in one thread (task.apply()) and asserts the second's lines do not carry the first's order_id.

Gotchas & Edge Cases

Prefork and context variables. Each prefork child process has its own context; signals run in the child that executes the task, so binding works. Do not bind in worker_process_init — that runs once per child, not per task.

gevent and eventlet pools. Greenlets share a thread; contextvars support depends on the library versions. Verify with a concurrency test, or bind context explicitly inside the task body for these pools.

Header propagation with canvas. Headers set on a signature propagate to the task they are set on, not automatically to every task in a chain; the before_task_publish handler covers this because each chain step is published from a context where request_id_var is set.

Result backend and retries. task.request.retries resets if a task is re-enqueued manually rather than with self.retry(). Use retry() so attempt numbers in logs stay truthful.

FAQ

Can I use Celery's built-in task logger instead? get_task_logger adds task name and id to its format, but only for records logged through it, and only as text. Structured context binding covers every logger and produces queryable fields.

Should the correlation id be the trace id? If you run OpenTelemetry, yes — propagate the trace context and log the trace id, as in propagating trace context through Celery tasks. A separate request_id is still useful when traces are sampled and logs are not.

What about tasks enqueued by Celery beat? Beat publishes from its own process, where no request is active, so request_id_var is empty and the before_task_publish handler adds nothing. The task_prerun fallback then uses the task's own id as the correlation id, which is the right root for a periodic run: every child task it enqueues inherits that id, and the whole run is searchable as one tree. Adding the schedule entry name as a header (headers={"schedule": "nightly_reconcile"}) makes periodic runs easy to filter as a group.

How much does this cost per task? A few microseconds for the signal handlers and context binding — negligible next to any real task.

Related