Instrumenting BullMQ with prom-client
BullMQ has no built-in Prometheus exporter, but it emits events for everything a worker does, and the prom-client library turns those events into metrics with a few dozen lines of code. This guide builds that instrumentation — throughput, failures, duration, wait time, and queue depth — as part of Prometheus Metrics for Workers in Observability & Monitoring for Job Queues.
Problem Statement
A Node.js service processes image thumbnails and outgoing emails with BullMQ. The only visibility is Bull Board, which shows the current state of each queue but no history. When users report slow thumbnails, nobody can say whether jobs are waiting longer, running longer, or failing and retrying. The platform team already runs Prometheus and Grafana for the API services. You want worker metrics in the same place, with alerting on backlog and failure rate, without adding a separate exporter service or high-cardinality labels that overload Prometheus.
Prerequisites
- BullMQ 5 workers running on Node.js 18 or newer.
prom-clientinstalled (npm install prom-client).- Prometheus able to scrape the worker pods or hosts.
- Agreement on a small set of labels: queue name and job name, not job IDs or user IDs.
Step 1 — Define the Metrics
Four metric types cover the questions teams ask about queues. Define them once in a module that every worker imports:
// metrics.ts
import client from "prom-client";
export const registry = new client.Registry();
client.collectDefaultMetrics({ register: registry }); // event loop lag, heap, CPU
export const jobsCompleted = new client.Counter({
name: "bullmq_jobs_completed_total",
help: "Jobs completed successfully",
labelNames: ["queue", "job_name"],
registers: [registry],
});
export const jobsFailed = new client.Counter({
name: "bullmq_jobs_failed_total",
help: "Job attempts that threw",
labelNames: ["queue", "job_name", "final"],
registers: [registry],
});
export const jobDuration = new client.Histogram({
name: "bullmq_job_duration_seconds",
help: "Processing time per attempt",
labelNames: ["queue", "job_name"],
buckets: [0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 300],
registers: [registry],
});
export const jobWait = new client.Histogram({
name: "bullmq_job_wait_seconds",
help: "Time from enqueue (or delay expiry) to processing start",
labelNames: ["queue", "job_name"],
buckets: [0.1, 0.5, 1, 5, 15, 60, 300, 900, 3600],
registers: [registry],
});
Choose histogram buckets from the durations you actually see; choosing histogram buckets for job duration covers how. The final label on failures distinguishes attempts that will be retried from jobs that have exhausted their attempts — the second is what users notice.
Step 2 — Record Metrics from Worker Events
The worker's completed and failed events carry the job, including its timestamps. processedOn is when the current attempt started, finishedOn when it ended, and timestamp when the job was created:
import { Worker, Job } from "bullmq";
import { jobsCompleted, jobsFailed, jobDuration, jobWait } from "./metrics";
export function instrument(worker: Worker) {
const queue = worker.name;
worker.on("active", (job: Job) => {
const readyAt = job.timestamp + (job.opts.delay ?? 0);
if (job.attemptsMade === 0 && job.processedOn) {
jobWait.labels(queue, job.name).observe(Math.max(0, job.processedOn - readyAt) / 1000);
}
});
worker.on("completed", (job: Job) => {
jobsCompleted.labels(queue, job.name).inc();
if (job.processedOn && job.finishedOn) {
jobDuration.labels(queue, job.name).observe((job.finishedOn - job.processedOn) / 1000);
}
});
worker.on("failed", (job: Job | undefined, err: Error) => {
if (!job) return;
const final = job.attemptsMade >= (job.opts.attempts ?? 1);
jobsFailed.labels(queue, job.name, String(final)).inc();
if (job.processedOn) {
jobDuration.labels(queue, job.name).observe((Date.now() - job.processedOn) / 1000);
}
});
}
Measuring wait from timestamp + delay excludes intended delays: a job scheduled for later should not count as waiting. Only the first attempt is recorded, because retries wait on purpose during backoff. The reasoning is covered in more depth in measuring queue wait time with enqueue timestamps.
Step 3 — Collect Queue Depth on Scrape
Counters and histograms describe work that happened. Backlog — how many jobs are waiting, delayed, or active right now — is state in Redis, so read it when Prometheus scrapes rather than on a timer. prom-client's collect callback runs at scrape time:
import { Queue } from "bullmq";
const queues = ["thumbnails", "emails"].map((n) => new Queue(n, { connection }));
new client.Gauge({
name: "bullmq_queue_jobs",
help: "Jobs in each state",
labelNames: ["queue", "state"],
registers: [registry],
async collect() {
this.reset();
for (const q of queues) {
const counts = await q.getJobCounts("waiting", "active", "delayed", "failed", "prioritized");
for (const [state, n] of Object.entries(counts)) this.labels(q.name, state).set(n);
}
},
});
Collect queue depth in exactly one place — a single worker replica or a small metrics sidecar — rather than in every worker process. Otherwise every replica reports the same numbers, and a sum() in Grafana multiplies the backlog by the replica count. Worker-level counters and histograms, in contrast, belong in every process, because each reports only its own work.
Step 4 — Expose the Endpoint
Serve the registry over HTTP from each worker process. A worker has no web framework, so Node's http module is enough:
import http from "node:http";
import { registry } from "./metrics";
http.createServer(async (req, res) => {
if (req.url !== "/metrics") { res.writeHead(404).end(); return; }
res.setHeader("Content-Type", registry.contentType);
res.end(await registry.metrics());
}).listen(Number(process.env.METRICS_PORT ?? 9464));
If a pod runs several Node processes with the cluster module, use prom-client's AggregatorRegistry so the primary process merges metrics from its children; otherwise each scrape sees only one child. With sandboxed processors, record metrics from the parent worker's events as above — the child process does not need to expose anything.
Add a Kubernetes PodMonitor or scrape annotations so Prometheus discovers the endpoint, and label targets with the service name so dashboards can filter. Keep the port separate from any HTTP port the worker might expose for health checks, so that a slow metrics collection never makes a liveness probe fail and restart a healthy worker.
Step 5 — Write the Queries and Alerts
With these metrics, the standard queue questions become short PromQL queries:
# throughput per queue (jobs/s)
sum by (queue) (rate(bullmq_jobs_completed_total[5m]))
# final failure ratio
sum by (queue) (rate(bullmq_jobs_failed_total{final="true"}[5m]))
/ sum by (queue) (rate(bullmq_jobs_completed_total[5m]) + rate(bullmq_jobs_failed_total{final="true"}[5m]))
# p95 wait time
histogram_quantile(0.95, sum by (queue, le) (rate(bullmq_job_wait_seconds_bucket[5m])))
# backlog
bullmq_queue_jobs{state="waiting"}
These four signals answer almost every first question during an incident. If throughput dropped but failures did not rise, workers are missing or stuck; if failures rose, look at the error logs for that job name; if wait time rose while throughput held steady, arrivals outgrew capacity and scaling is the fix; if backlog grows steadily over hours, capacity is short even when nothing is failing.
Alert on wait time and on final failures rather than on raw backlog alone — a large backlog that drains within its target is fine. Alerting on queue backlog with Prometheus has complete rules, and building a BullMQ Grafana dashboard turns these queries into panels.
Verification
curl localhost:9464/metricson a worker shows the counters and histograms withqueueandjob_namelabels, and they increase as jobs run.- The depth gauge appears from one target only, and matches Bull Board's counts.
- A job made to throw with
attempts: 3produces twofinal="false"failures and onefinal="true". - A delayed job does not appear as a long wait in
bullmq_job_wait_seconds. - The total number of series from these metrics stays in the hundreds, not thousands.
Gotchas & Edge Cases
Unbounded label values. Job names built from data (resize-${userId}) create a series per user. Keep job names to a fixed set and put variable data in the payload.
Events and process restarts. Counters reset when a process restarts; Prometheus's rate() handles resets, but raw counter values in dashboards will drop. Always graph rates.
Clock skew. timestamp is set by the producer's clock and processedOn by the worker's. Skew between machines shows up as wait time; keep NTP running, and treat negative values as zero.
Event-loop lag hides in the defaults. collectDefaultMetrics exports nodejs_eventloop_lag_seconds. A worker with high lag fetches jobs slowly and renews locks late, which shows up as rising wait time and stalled jobs. Put it on the dashboard next to the job metrics; it is often the first place a CPU-heavy processor shows itself.
Stalled jobs. A job whose lock expires is moved back to waiting without a failed event. Listen for stalled on a QueueEvents instance and count it separately, since stalls usually indicate event-loop blocking or crashes.
FAQ
Should I use an existing BullMQ exporter instead? Community exporters read queue depth from Redis and are useful for that. They cannot see per-attempt duration or wait time as precisely as the worker itself, so most teams combine an exporter (or the gauge above) with worker-side instrumentation.
Do I need OpenTelemetry metrics instead of prom-client? Either works with Prometheus. prom-client is simpler if you only need Prometheus; OpenTelemetry is the better choice if you also send traces, as in OpenTelemetry tracing for BullMQ.
How much overhead does this add? Observing a histogram and incrementing a counter take microseconds. The depth gauge costs one Redis round trip per queue per scrape, which is negligible at a 15–30 second scrape interval.
Can I label metrics by tenant? Only if the number of tenants is small and fixed. For thousands of tenants, record tenant-level detail in logs or traces and keep metrics at queue and job-name level.
Related
- Prometheus Metrics for Workers — what to measure and why.
- Instrumenting Celery with a Prometheus Exporter — the Python counterpart.
- Building a BullMQ Grafana Dashboard — visualising these metrics.
- Measuring Queue Wait Time with Enqueue Timestamps — the wait-time metric in depth.