事件、进度与任务结果

任务进入等待、开始、完成或失败时,BullMQ 会发出事件。事件适合日志、监控、WebSocket 通知和管理后台;任务进度适合展示长任务阶段;返回值适合小型结果,但关键业务结果应可靠地写入数据库或结果队列。

Worker 本地事件

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

import { Job, Worker } from "bullmq";

interface ExportJobData {
  exportId: string;
}

interface ExportResult {
  objectKey: string;
  rowCount: number;
}

const worker = new Worker<ExportJobData, ExportResult>(
  "export",
  async (job): Promise<ExportResult> => {
    return buildExport(job);
  },
  { connection: { host: "127.0.0.1", port: 6379 } },
);

worker.on("completed", (job: Job<ExportJobData, ExportResult>, result) => {
  console.log("导出完成", { jobId: job.id, result });
});

worker.on("failed", (job, error: Error) => {
  console.error("导出失败", { jobId: job?.id, message: error.message });
});

worker.on("error", (error: Error) => {
  console.error("Worker 自身错误", error);
});

async function buildExport(
  job: Job<ExportJobData, ExportResult>,
): Promise<ExportResult> {
  console.log(`处理导出 ${job.data.exportId}`);
  return { objectKey: "exports/report.csv", rowCount: 100 };
}

failed 是某个任务失败;error 是 Worker 自身发生连接等错误。生产代码应同时监听二者。

QueueEvents 全局事件

多个 Worker 或多台机器共同处理队列时,用 QueueEvents 在一个地方观察整个队列:

import { QueueEvents } from "bullmq";

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

await queueEvents.waitUntilReady();

queueEvents.on("completed", ({ jobId, returnvalue }) => {
  console.log("任意 Worker 完成任务", { jobId, returnvalue });
});

queueEvents.on("failed", ({ jobId, failedReason }) => {
  console.error("任意 Worker 执行失败", { jobId, failedReason });
});

queueEvents.on("progress", ({ jobId, data }) => {
  console.log("任务进度变化", { jobId, data });
});

QueueEvents 使用 Redis Streams,而非普通 Pub/Sub,因此短暂断线时具有更好的事件交付特性。但事件流会自动裁剪,默认大约保留 10,000 条;它不是永久审计日志。

更新任务进度

interface ImportProgress {
  stage: "reading" | "validating" | "saving";
  percent: number;
}

const importWorker = new Worker<ExportJobData>(
  "import",
  async (job): Promise<void> => {
    const reading: ImportProgress = { stage: "reading", percent: 10 };
    await job.updateProgress(reading);

    const validating: ImportProgress = {
      stage: "validating",
      percent: 50,
    };
    await job.updateProgress(validating);

    const saving: ImportProgress = { stage: "saving", percent: 90 };
    await job.updateProgress(saving);
  },
  { connection: { host: "127.0.0.1", port: 6379 } },
);

进度既可以是数字,也可以是可序列化对象。对象比单一百分比更能表达当前阶段,但不要包含大块数据。

实用建议:

返回值与 returnvalue

Worker 处理函数的返回值会保存在任务的 returnvalue 中:

const resultWorker = new Worker<ExportJobData, ExportResult>(
  "export-result",
  async (job): Promise<ExportResult> => {
    return {
      objectKey: `exports/${job.data.exportId}.csv`,
      rowCount: 250,
    };
  },
  { connection: { host: "127.0.0.1", port: 6379 } },
);

返回值应小而可序列化。不要把 CSV、图片或大 JSON 直接放入 Redis;将文件存入对象存储,只返回 objectKey、URL 标识或数据库记录 ID。

如何可靠保存关键结果

只在 completed 事件监听器中写数据库有风险:任务已经标记完成后,监听器可能崩溃,最终出现“任务成功但结果没保存”。官方建议更可靠的方式是:

  1. 在处理函数返回前保存业务结果,保存成功后任务才完成。
  2. 或将结果投递到专门的 results 队列,再由结果 Worker 持久化。
const reliableWorker = new Worker<ExportJobData, ExportResult>(
  "export",
  async (job): Promise<ExportResult> => {
    const result = await createExport(job.data.exportId);
    await saveExportResult(job.data.exportId, result);
    return result;
  },
  { connection: { host: "127.0.0.1", port: 6379 } },
);

保存操作也要幂等,例如以 exportId 建立唯一约束。否则任务重试时可能创建多条结果记录。

事件保留与清理

可通过队列配置调整事件流长度,也可手动保留最近若干事件:

import { Queue } from "bullmq";

const exportQueue = new Queue("export", {
  connection: { host: "127.0.0.1", port: 6379 },
  streams: {
    events: { maxLen: 20_000 },
  },
});

await exportQueue.trimEvents(10_000);

增大事件流会增加 Redis 内存占用。长期审计应写入日志平台或数据库,而不是无限保留 QueueEvents。

最佳实践清单

本章小结

Worker 事件是实例本地视角,QueueEvents 是跨 Worker 的队列视角;updateProgress() 用于可观察进度,处理函数返回值进入 returnvalue。关键结果必须在任务完成前可靠保存,不能把事件监听器当作事务提交点。

官方文档