Structured JSON Logging for BullMQ Workers

Node workers running BullMQ often log with console.log or a single global logger, which leaves every line without the job that produced it. This guide sets up pino so every line carries job context automatically, as part of Structured Logging for Workers in Observability & Monitoring for Job Queues. It covers child loggers per job, AsyncLocalStorage for code that cannot receive a logger argument, attempt-aware levels for worker events, redaction, and sandboxed processors.

Problem Statement

A TypeScript service processes document-conversion jobs with BullMQ at a concurrency of 50 per worker. Its logs are console.log strings interleaved from 50 concurrent jobs, so a line like Converting page 14 cannot be attributed to any document. Failed jobs appear only as Job failed from a global worker.on("failed") handler, logged at error for every attempt — 30,000 error lines a day, almost all from jobs that later succeeded. You want every line tied to its job id and document id, error-level lines only for final failures, no personal data in logs, and consistent output from sandboxed processors running in child processes.

Prerequisites

  • Node 20+, BullMQ 5.x, and pino 9 (plus pino-pretty for local development only).
  • Processors written as functions receiving the Job object.
  • A log collector reading container stdout (Promtail, Alloy, Fluent Bit, Vector, or a vendor agent).
  • Optionally OpenTelemetry, if you want trace ids on log lines (Step 6).

Step 1 — Create a Base Logger with a Fixed Schema

Configure pino once with the service's static fields, ISO timestamps, level labels instead of numbers, and redaction paths for fields that must never be logged.

// src/log.ts
import pino from "pino";

export const baseLogger = pino({
  level: process.env.LOG_LEVEL ?? "info",
  base: { service: "doc-converter", env: process.env.DEPLOY_ENV },
  timestamp: pino.stdTimeFunctions.isoTime,
  formatters: { level: (label) => ({ level: label }) },       // "error", not 50
  messageKey: "msg",
  redact: {
    paths: ["data.email", "data.token", "*.password", "req.headers.authorization"],
    censor: "[redacted]",
  },
  serializers: { err: pino.stdSerializers.err },             // stack as a field, one line
});

pino writes one JSON object per line to stdout synchronously by default in production-safe mode; avoid transports that run in worker threads for the main stream unless you flush them on shutdown, or you lose the last lines before a crash.

Base logger, job child, step child The base logger carries service and environment. When a job starts, a child logger adds job id, job name, queue, attempt, and document id. Inside the processor, a further child can add a step name such as render page. Every line from any level of the hierarchy includes all of its ancestors' fields. Fields accumulate down the hierarchy baseLogger service, env job child job_id, attempt, document_id step child step: render_page Child creation is cheap: pino pre-serializes bound fields once per child, not per line.

Step 2 — Create a Child Logger per Job

Wrap every processor so it receives a logger already bound to the job. The wrapper is the single place where job fields are chosen, which keeps the schema consistent across processors.

// src/with-job-logger.ts
import type { Job } from "bullmq";
import type { Logger } from "pino";
import { baseLogger } from "./log";

const BUSINESS_KEYS = ["documentId", "tenantId", "userId"] as const;

export function jobLogger(job: Job): Logger {
  const ids: Record<string, unknown> = {};
  for (const k of BUSINESS_KEYS) if (job.data?.[k] !== undefined) ids[k] = job.data[k];
  return baseLogger.child({
    job_id: job.id,
    job_name: job.name,
    queue: job.queueName,
    attempt: job.attemptsMade + 1,                   // attemptsMade counts completed attempts
    max_attempts: job.opts.attempts ?? 1,
    enqueued_at: new Date(job.timestamp).toISOString(),
    ...ids,
  });
}

export function withJobLogger<T, R>(fn: (job: Job<T>, log: Logger) => Promise<R>) {
  return async (job: Job<T>): Promise<R> => {
    const log = jobLogger(job);
    const started = Date.now();
    log.debug("job started");
    const result = await fn(job, log);
    log.info({ duration_ms: Date.now() - started }, "job succeeded");
    return result;
  };
}
// src/processors/convert.ts
export const convert = withJobLogger<ConvertData, ConvertResult>(async (job, log) => {
  const doc = await loadDocument(job.data.documentId);
  for (const [i, page] of doc.pages.entries()) {
    log.debug({ page: i + 1 }, "rendering page");
    await renderPage(page);
    await job.updateProgress(Math.round(((i + 1) / doc.pages.length) * 100));
  }
  return { pages: doc.pages.length };
});

Business ids come from a whitelist, never from spreading job.data, which may contain file contents or personal data.

Step 3 — Use AsyncLocalStorage for Code That Cannot Take a Logger

Deep library code — a storage client, a PDF renderer wrapper — should not need a log parameter threaded through every call. AsyncLocalStorage makes the job's logger available anywhere in the async call tree started by the processor.

// src/context.ts
import { AsyncLocalStorage } from "node:async_hooks";
import type { Logger } from "pino";
import { baseLogger } from "./log";

const als = new AsyncLocalStorage<Logger>();

export const runWithLogger = <R>(log: Logger, fn: () => Promise<R>) => als.run(log, fn);
export const currentLogger = (): Logger => als.getStore() ?? baseLogger;

// in withJobLogger: const result = await runWithLogger(log, () => fn(job, log));

// src/storage.ts — no logger parameter needed
export async function putObject(key: string, body: Buffer) {
  const log = currentLogger();                       // the calling job's logger
  const t0 = Date.now();
  await s3.send(new PutObjectCommand({ Bucket: BUCKET, Key: key, Body: body }));
  log.debug({ key, bytes: body.length, ms: Date.now() - t0 }, "stored object");
}

With 50 concurrent jobs in one process, each async chain sees its own store, so lines never cross-attribute — the Node equivalent of the context-leak protection described in Structured Logging for Workers.

Concurrent jobs, isolated context Job 17 and job 42 run concurrently in the same event loop. Each processor call runs inside its own AsyncLocalStorage store holding that job's child logger. When the shared storage client logs, it reads the store for the current async chain, so job 17's upload line carries job 17's ids and job 42's carries job 42's. Same function, different job context job 17: store = log17 job 42: store = log42 putObject() currentLogger() line: job_id=17 line: job_id=42

Step 4 — Log Worker Events with Attempt-Aware Levels

BullMQ's failed event fires on every failed attempt. Decide the level from attemptsMade versus the configured attempts, so only jobs that have given up log at error.

// src/worker.ts
import { Worker, UnrecoverableError } from "bullmq";

const worker = new Worker("convert", convert, { connection, concurrency: 50 });

worker.on("failed", (job, err) => {
  if (!job) return;
  const log = jobLogger(job).child({ attempt: job.attemptsMade });  // attemptsMade now includes this one
  const max = job.opts.attempts ?? 1;
  const final = job.attemptsMade >= max || err instanceof UnrecoverableError;
  const fields = { err, error_class: err.name, final };
  if (final) log.error(fields, "job failed permanently");
  else if (job.attemptsMade === max - 1) log.warn(fields, "job failed, last retry scheduled");
  else log.info(fields, "job failed, will retry");
});

worker.on("stalled", (jobId) => baseLogger.warn({ job_id: jobId, queue: "convert" }, "job stalled"));
worker.on("error", (err) => baseLogger.error({ err }, "worker error"));   // connection-level

The stalled event deserves its own line and alert: a stalled job means a worker lost its lock, usually because the event loop was blocked — see BullMQ lock duration and stalled jobs. With these levels, the 30,000 daily error lines from the problem statement drop to the few dozen jobs that actually failed.

That reduction is what makes error-level logs usable as an alert source again. An alert rule on "any level=error line from doc-converter in the last five minutes" was useless when retries produced a steady stream of them; after the change, each error line is a document that will not be converted without intervention, and the alert carries the document id in its payload.

Error lines per day Before the change, every failed attempt logged at error, producing about thirty thousand error lines a day, most from jobs that succeeded on a later attempt. After the change, retries log at info or warning and only final failures log at error, about forty lines a day, each one a document that needs attention. level="error" lines per day every attempt ~30,000: mostly jobs that later succeeded final only ~40: each one a document that needs a human Retries still appear in logs, at info and warning, with the same job context.

Step 5 — Keep Sandboxed Processors Consistent

Sandboxed processors run in child processes (see BullMQ sandboxed processors). They cannot share the parent's logger instance, but they write to the same stdout, so create the logger inside the child with the same configuration module.

// src/processors/convert.sandbox.ts — loaded by the Worker as a file path
import type { SandboxedJob } from "bullmq";
import { baseLogger } from "../log";                  // same config, new instance in the child

export default async function (job: SandboxedJob) {
  const log = baseLogger.child({ job_id: job.id, job_name: job.name, queue: job.queueName,
                                 attempt: job.attemptsMade + 1, documentId: job.data.documentId,
                                 sandbox_pid: process.pid });
  log.info("job started in sandbox");
  // ... heavy CPU work ...
  return { ok: true };
}

Adding sandbox_pid helps when a child process crashes: the last lines from that pid show what it was doing.

Step 6 — Add Trace Ids for Log-to-Trace Links

If OpenTelemetry instruments BullMQ (as in OpenTelemetry tracing for BullMQ), add the active span's ids through a pino mixin so every line links to its trace.

import { trace } from "@opentelemetry/api";

export const baseLogger = pino({
  // ...options from Step 1...
  mixin() {
    const ctx = trace.getActiveSpan()?.spanContext();
    return ctx ? { trace_id: ctx.traceId, span_id: ctx.spanId } : {};
  },
});

Grafana and most vendors turn a trace_id field into a clickable link to the trace, closing the loop from "this job failed" to "this is where its time went".

Verification

# Local: run a job and inspect lines
LOG_LEVEL=debug node dist/worker.js | npx pino-pretty
# Staging: every line for one document, across attempts
logcli query '{service="doc-converter"} | json | documentId="d-5521"' --limit 100
# Error volume should now track true final failures
logcli query 'sum(count_over_time({service="doc-converter"} | json | level="error" [24h]))'

Add a test that runs two jobs concurrently with a shared helper that logs, captures pino output with a destination stream, and asserts each line's job_id matches the job whose chain produced it.

Gotchas & Edge Cases

attemptsMade semantics. Inside the processor, attemptsMade is the number of completed attempts (so the current attempt is attemptsMade + 1); in the failed handler it already includes the failed one. Getting this wrong makes every final failure look like a retry.

Logging whole jobs. log.info({ job }) serializes the entire job, including data and return values. Use the child logger's bound fields instead.

Asynchronous transports and crashes. pino transports in worker threads buffer; a crash can drop the final lines. Keep stdout as the primary destination and let the platform ship logs.

console.log in dependencies. Libraries that log with console bypass pino. Redirect console.log to currentLogger() in the worker entry point if they matter.

FAQ

Winston or pino? Either works with the same pattern (child loggers plus AsyncLocalStorage). pino is substantially faster at high line rates, which matters for workers at concurrency 50+.

Should updateProgress also log? No — progress is state, not an event. Log at step boundaries, and let progress live in the job record.

How do I log from QueueEvents listeners in another service? Create a logger with the job id from the event payload; you do not have the job object, so fetch it with Job.fromId only if you need more context.

Related