16 - 指标、事件与链路追踪

队列系统不能只监控“进程还活着吗”。生产环境真正关心的是:积压是否增长、任务是否变慢、失败是否突然增加,以及任务时间究竟花在哪个服务。

BullMQ 提供三组互补能力:

1. 启用 BullMQ 内置 Metrics

Metrics 在 Worker 上启用。每个数据点代表一分钟内完成或失败的任务数,数据保存在 Redis list 中,超过保留上限后自动清理。

import { Job, MetricsTime, Worker } from 'bullmq';

interface PaintJobData {
  color: string;
}

const worker = new Worker<PaintJobData>(
  'paint',
  async (job: Job<PaintJobData>): Promise<void> => {
    console.log(`使用 ${job.data.color} 开始喷涂`);
  },
  {
    connection: { host: '127.0.0.1', port: 6379 },
    metrics: {
      maxDataPoints: MetricsTime.ONE_WEEK * 2,
    },
  },
);

官方特别提醒:同一队列的所有 Worker 应使用相同的 metrics 配置,否则统计结果可能不一致。官方估算,保存两周的一分钟粒度 completed/failed 数据,每个队列大约只占 120 KB Redis 内存。

2. 查询内置 Metrics

import { MetricsTime, Queue } from 'bullmq';

const queue = new Queue('paint', {
  connection: { host: '127.0.0.1', port: 6379 },
});

const completed = await queue.getMetrics(
  'completed',
  0,
  MetricsTime.ONE_WEEK * 2,
);

const failed = await queue.getMetrics(
  'failed',
  0,
  MetricsTime.ONE_WEEK * 2,
);

console.log({
  completedTotal: completed.meta.count,
  completedPerMinute: completed.data,
  failedTotal: failed.meta.count,
});

data 数组的每个位置代表一分钟。meta.count 是队列开始处理以来的累计数,不只是当前查询区间的合计。startend 参数可用于分页。

内置 Metrics 很轻量,但它不直接提供端到端延迟、任务数据大小、业务成功率等全部指标,这些仍需应用自行补充。

3. 导出 Prometheus 格式

Queue 的 exportPrometheusMetrics() 可以输出 Prometheus 文本格式。下面使用 Node.js 内置 HTTP 模块,不引入 Web 框架:

import http, { IncomingMessage, ServerResponse } from 'node:http';
import { Queue } from 'bullmq';

const queue = new Queue('paint', {
  connection: { host: '127.0.0.1', port: 6379 },
});

const server = http.createServer(
  async (request: IncomingMessage, response: ServerResponse): Promise<void> => {
    if (request.method !== 'GET' || request.url !== '/metrics') {
      response.writeHead(404).end('Not Found');
      return;
    }

    try {
      const metrics = await queue.exportPrometheusMetrics({
        environment: 'development',
        service: 'paint-worker',
      });

      response.writeHead(200, {
        'Content-Type': 'text/plain; version=0.0.4',
      });
      response.end(metrics);
    } catch (error: unknown) {
      const message = error instanceof Error ? error.message : 'Unknown error';
      response.writeHead(500).end(message);
    }
  },
);

const port = 9464;
server.listen(port, () => {
  console.log(`Metrics: http://127.0.0.1:${port}/metrics`);
});

附加标签适合区分环境和服务,但不要把 jobId、用户 ID 等高基数字段作为 Prometheus label,否则会产生大量时间序列。

4. Worker 本地事件与 QueueEvents

Worker 事件只反映这个 Worker 实例处理到的任务:

worker.on('completed', (job: Job<PaintJobData>) => {
  console.log(`本 Worker 完成任务 ${job.id}`);
});

worker.on('failed', (job: Job<PaintJobData> | undefined, error: Error) => {
  console.error(`任务 ${job?.id ?? 'unknown'} 失败`, error);
});

worker.on('error', (error: Error) => {
  console.error('Worker 基础设施错误', error);
});

如果要在一个地方监听所有 Worker 的事件,应使用 QueueEvents

import { QueueEvents } from 'bullmq';

const queueEvents = new QueueEvents('paint', {
  connection: { host: '127.0.0.1', port: 6379 },
});

queueEvents.on('completed', ({ jobId }: { jobId: string }) => {
  console.log(`任意 Worker 完成任务 ${jobId}`);
});

queueEvents.on(
  'progress',
  ({ jobId, data }: { jobId: string; data: number | object }) => {
    console.log(`任务 ${jobId} 进度`, data);
  },
);

queueEvents.on('error', (error: Error) => {
  console.error('QueueEvents 错误', error);
});

QueueEvents 使用 Redis Streams,而不是普通 Pub/Sub,因此断线场景下具有更好的事件交付特性。事件流会自动裁剪,默认大约保留 10,000 条;可通过 streams.events.maxLen 调整,也可调用 queue.trimEvents() 手动裁剪。

事件监听适合触发日志、WebSocket 通知或观测逻辑,但关键业务一致性仍应由数据库状态、幂等约束或 Flow 保证,不能假设监听器绝不会重复或漏处理。

5. Telemetry 与 OpenTelemetry 概览

BullMQ 提供 Telemetry 接口,官方当前提供 bullmq-otel 来接入 OpenTelemetry。它适合回答:

官方入门方式需要额外安装可选包:

npm install bullmq-otel

本教程不会自动安装依赖。使用时可把同一个 telemetry 实现传给 Queue 和 Worker:

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

const telemetry = new BullMQOtel('report-service');
const connection = { host: '127.0.0.1', port: 6379 };

const queue = new Queue('reports', { connection, telemetry });

const worker = new Worker(
  'reports',
  async (): Promise<string> => 'done',
  {
    name: 'report-worker',
    connection,
    telemetry,
  },
);

OpenTelemetry 还需要 exporter、collector 或 Jaeger 等后端配置。bullmq-otel 让 BullMQ 产生 spans,不代表完整可观测平台已经配置完成。

6. 建议监控的核心指标

类别 指标 说明
积压 waiting、delayed 数量 持续增长说明生产速度超过消费速度
处理中 active 数量 长期接近并发上限可能需要扩容或优化
结果 completed、failed 每分钟数量 观察吞吐和失败突增
延迟 入队到开始、开始到完成的时间 区分排队慢和处理慢
可靠性 stalled、重试次数、lock 错误 暴露 CPU 阻塞、网络或锁续订问题
Redis 内存、延迟、连接数、持久化状态 Redis 是队列正确运行的基础

告警应关注趋势,而不是只设固定数量。例如等待任务为 1,000 不一定异常,但连续 15 分钟只增不减通常值得告警。

7. 最佳实践清单

官方资料