16 - 指标、事件与链路追踪
队列系统不能只监控“进程还活着吗”。生产环境真正关心的是:积压是否增长、任务是否变慢、失败是否突然增加,以及任务时间究竟花在哪个服务。
BullMQ 提供三组互补能力:
- Metrics:按分钟统计 completed/failed 数量。
- Events:观察任务生命周期事件。
- Telemetry:把任务生命周期接入 OpenTelemetry 链路追踪。
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 是队列开始处理以来的累计数,不只是当前查询区间的合计。start 和 end 参数可用于分页。
内置 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。它适合回答:
- 哪个服务创建了任务?
- 任务排队等待了多久?
- Worker 处理花了多久?
- 任务调用的数据库或 HTTP 服务哪里最慢?
官方入门方式需要额外安装可选包:
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. 最佳实践清单
- Metrics、日志和 traces 使用同一个 queue name、job name、job ID 进行关联。
- 不在日志或 telemetry 中记录完整敏感 job data。
- 同一队列的所有 Worker 使用一致的 metrics 配置。
failed、stalled、error分开统计,它们代表不同问题。- Prometheus label 保持低基数。
- 为
/metrics端点设置适当的网络访问控制。 - 应用关闭时依次关闭 Worker、QueueEvents、Queue 和监控 HTTP Server。