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.
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.
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.
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
- Distributed Tracing for Async Jobs — why async tracing is different.
- Propagating Trace Context Through Celery Tasks — the Python equivalent.
- Tracing Sidekiq Jobs with OpenTelemetry — the Ruby equivalent.
- Instrumenting BullMQ with prom-client — metrics alongside traces.