OpenTelemetry Tracing for BullMQ

Tracing an HTTP request is straightforward because the request and response happen in one call chain. A BullMQ job breaks that chain: the request ends when the job is added, and the work happens later in a different process. OpenTelemetry can join the two if the trace context travels with the job. This guide sets that up for BullMQ, as part of Distributed Tracing for Async Jobs in Observability & Monitoring for Job Queues.

Problem Statement

A Node.js API accepts video uploads, stores them, and adds a BullMQ job to transcode them. Users report that some uploads take ten minutes to become available. The API's traces end in 80 ms at the queue.add call; the worker's traces, when they exist, are separate and have no connection to the upload request. Nobody can tell whether the time is spent waiting in the queue, transcoding, or retrying a failed step. You want a single trace that shows the upload request, the time the job waited, the worker's processing with its downstream calls, and any retries, so slow uploads can be explained from one view.

Prerequisites

  • BullMQ 5.x and Node.js 18 or newer on both the API and the workers.
  • The OpenTelemetry Node SDK (@opentelemetry/sdk-node) with auto-instrumentation for HTTP and your database clients.
  • A tracing backend reachable through OTLP: Jaeger, Tempo, Honeycomb, or similar, usually via an OpenTelemetry Collector.
  • The same service naming convention on both sides (video-api, video-worker).

Step 1 — Initialise OpenTelemetry in Both Processes

The SDK must start before any instrumented module is loaded. Create a small tracing.ts and load it first in both the API and the worker (node --import ./tracing.js or -r for CommonJS):

// tracing.ts
import { NodeSDK } from "@opentelemetry/sdk-node";
import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http";
import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node";

const sdk = new NodeSDK({
  serviceName: process.env.OTEL_SERVICE_NAME,           // video-api or video-worker
  traceExporter: new OTLPTraceExporter({ url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT + "/v1/traces" }),
  instrumentations: [getNodeAutoInstrumentations({ "@opentelemetry/instrumentation-fs": { enabled: false } })],
});
sdk.start();
process.on("SIGTERM", () => sdk.shutdown());

Auto-instrumentation traces HTTP, database, and Redis calls, but it does not know that a Redis write from BullMQ is "adding a job" or that a later read is "processing" it. That link is what the next steps add.

Step 2 — Enable BullMQ's Built-in Telemetry

BullMQ 5 has a telemetry option that creates producer and consumer spans and carries trace context inside the job. Install the OpenTelemetry adapter and pass it to the Queue and the Worker:

import { Queue, Worker } from "bullmq";
import { BullMQOtel } from "bullmq-otel";

const telemetry = new BullMQOtel("video-pipeline");

export const transcodeQueue = new Queue("transcode", { connection, telemetry });

new Worker("transcode", processTranscode, { connection, telemetry, concurrency: 4 });

With telemetry enabled, queue.add creates a producer span as a child of the active span (the HTTP request), serialises its context into the job's metadata, and the worker creates a consumer span for each attempt linked to that context. Retries produce a new consumer span each time, all attached to the same trace.

A single trace from upload to transcoded video The trace starts with the HTTP POST upload span of 80 milliseconds, containing a producer span for adding the transcode job. The job then waits in the queue, visible as a gap. The worker's first consumer span starts later, with child spans for downloading the source, running the transcoder and storing the output; it fails, and after backoff a second consumer span succeeds. All spans share one trace ID because the context travelled with the job. One trace ID, two services, a gap in between video-api POST /uploads transcode add (producer) waiting in queue video-worker process attempt 1 download, transcode (error) process attempt 2 download, transcode, store Queue wait and retry backoff appear as gaps you can measure.

Step 3 — Propagate Context Manually If You Cannot Use the Option

On older BullMQ versions, or with a custom queue wrapper, carry the context yourself with the W3C Trace Context format. Inject it into the job data on the producer side and extract it in the worker:

import { context, propagation, trace, SpanKind, ROOT_CONTEXT } from "@opentelemetry/api";
const tracer = trace.getTracer("video-pipeline");

export async function addTraced(name: string, data: object, opts = {}) {
  return tracer.startActiveSpan(`${name} publish`, { kind: SpanKind.PRODUCER }, async (span) => {
    const carrier: Record<string, string> = {};
    propagation.inject(context.active(), carrier);          // traceparent, tracestate
    try { return await transcodeQueue.add(name, { ...data, _otel: carrier }, opts); }
    finally { span.end(); }
  });
}

export function traced(processor: (job: Job) => Promise<unknown>) {
  return async (job: Job) => {
    const parent = propagation.extract(ROOT_CONTEXT, job.data._otel ?? {});
    return tracer.startActiveSpan(`${job.name} process`,
      { kind: SpanKind.CONSUMER, attributes: { "messaging.system": "bullmq", "messaging.message.id": job.id } },
      parent,
      async (span) => {
        try { return await processor(job); }
        catch (e) { span.recordException(e as Error); span.setStatus({ code: 2 }); throw e; }
        finally { span.end(); }
      });
  };
}

Wrap the processor (new Worker("transcode", traced(processTranscode), ...)) and every span created inside it — HTTP calls, S3 uploads, database queries — becomes part of the upload's trace. The same approach for Python is shown in propagating trace context through Celery tasks.

Step 4 — Add Attributes That Answer Queue Questions

Default spans show timing but not queue-specific facts. Add attributes to the consumer span that you will want to filter and group by:

span.setAttributes({
  "messaging.destination.name": job.queueName,
  "bullmq.job.attempt": job.attemptsMade + 1,
  "bullmq.job.wait_ms": (job.processedOn ?? Date.now()) - job.timestamp - (job.opts.delay ?? 0),
  "video.size_mb": job.data.sizeMb,
});

wait_ms on the span lets you query "slowest waits in the last hour" in the tracing backend and click straight into examples. Keep attribute values small and free of personal data; traces are often retained and searched more widely than logs.

Step 5 — Sample Sensibly Across the Queue

A trace that is sampled on the API side must also be sampled in the worker, or it will have a hole. Parent-based sampling handles this: the worker respects the sampling decision carried in traceparent.

Keeping sampling decisions consistent The API makes the sampling decision once, at the root span, for example 10 percent of requests. The sampled flag is part of the traceparent value stored in the job. The worker uses a parent-based sampler, so it records spans only for jobs whose flag says sampled. Complete traces are kept and unsampled ones are dropped on both sides, avoiding traces with missing halves. Decide once at the root, follow everywhere API root span sample 10% job metadata traceparent …-01 (sampled) worker ParentBased sampler Jobs enqueued without a trace (cron, scripts) start a new root in the worker. Tail sampling in the Collector can keep all slow or failed traces.

Set OTEL_TRACES_SAMPLER=parentbased_traceidratio and OTEL_TRACES_SAMPLER_ARG=0.1 on both services. Jobs added by scheduled job schedulers or scripts have no parent, so the worker makes its own decision for them. To keep every slow or failed job regardless of the ratio, add tail-based sampling in the OpenTelemetry Collector with policies for latency and error status.

Step 6 — Read the Trace

With everything connected, open a slow upload's trace. The gap between the producer span and the first consumer span is queue wait — if it is large, the fix is more workers or better concurrency settings. Several consumer spans mean retries, and the failed ones carry the exception. A long single consumer span with one dominant child points to the slow step. In the video example, traces showed most slow uploads had three failed attempts against a storage endpoint with a short timeout, and the fix was a timeout change rather than more workers.

Reading the shape of a job trace Shape one: a short producer span, then a long empty gap, then a normal consumer span; the job waited in the queue, so capacity or concurrency is the issue. Shape two: several consumer spans separated by backoff gaps, the earlier ones with errors; retries against a failing dependency. Shape three: a single long consumer span where one child span takes most of the time; that step is slow. Three shapes, three different fixes long wait → add capacity retries → fix the dependency slow step → optimise that step

Most tracing backends can search by these shapes directly. Query for consumer spans where bullmq.job.wait_ms exceeds your target to find waiting problems; for traces with more than one consumer span, or spans with bullmq.job.attempt above 1, to find retry problems; and for consumer spans above a duration threshold, grouped by their slowest child span name, to find slow steps. Saving these three searches next to the dashboard turns a vague "uploads are slow" report into a specific cause within minutes. Pair them with the aggregate view in metrics: if the traces show waiting, the wait-time histogram should confirm it across all jobs rather than a handful of examples.

Verification

  • A test upload produces one trace containing the HTTP span, the producer span, and the worker's consumer span, all with the same trace ID.
  • A job that fails once and then succeeds shows two consumer spans, the first with error status and the exception.
  • With a 10% sampling ratio, sampled traces are complete; no worker-only fragments exist for sampled API requests.
  • Jobs added from a cron script produce their own root traces in the worker.

Gotchas & Edge Cases

Loading order. If BullMQ or ioredis is imported before the SDK starts, auto-instrumentation misses them. Load the tracing file with --import or -r.

Very long traces. A job that waits hours makes a trace that spans hours; some backends have maximum trace durations or search windows. Consider span links instead of parent-child for long waits — see span links for batch and fan-out jobs.

Sandboxed processors. With sandboxed processors, the job runs in a child process; initialise the SDK in the child as well, or spans from inside the processor are lost.

Flows. In BullMQ flows, children are added with the parent's context. The trace shows the whole tree, which is useful, but large flows produce very wide traces.

FAQ

Does tracing slow jobs down? Creating spans costs microseconds. Exporting is batched in the background. The overhead is negligible compared with typical job work.

Can I trace without a Collector? Yes, the SDK can export directly to a backend. A Collector adds batching, tail sampling, and a single place to change destinations, which is worth it beyond the smallest setups.

How do traces relate to metrics? Metrics tell you that wait time rose; traces show why for specific jobs. Link them with exemplars, as in Grafana heatmaps for job duration.

What about jobs that never run? A job that is removed, expires, or sits in the failed set after its last attempt still has its producer span, but the consumer spans stop at the final failure. Searching for producer spans with no successful consumer is a quick way to find work that users are still waiting on.

Related