事件、进度与任务结果
任务进入等待、开始、完成或失败时,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 } },
);
进度既可以是数字,也可以是可序列化对象。对象比单一百分比更能表达当前阶段,但不要包含大块数据。
实用建议:
- 不要在每条记录后都更新进度,会制造大量 Redis 写入和事件。
- 按阶段或每处理固定批次更新,例如每 1000 条一次。
- 百分比只是展示值,任务最终状态仍以 completed / failed 为准。
- 前端断线重连后应主动查询任务状态,不能只依赖实时事件补齐界面。
返回值与 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 事件监听器中写数据库有风险:任务已经标记完成后,监听器可能崩溃,最终出现“任务成功但结果没保存”。官方建议更可靠的方式是:
- 在处理函数返回前保存业务结果,保存成功后任务才完成。
- 或将结果投递到专门的 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 监听
error,任务失败监听failed。 - 跨进程观察使用
QueueEvents,局部处理使用 Worker 事件。 - 事件处理器保持轻量;耗时操作进入另一条队列。
- 关键结果在任务处理函数内持久化,并保证幂等。
returnvalue只保存小型元数据,不保存大文件。- 进度更新节流,前端同时支持主动查询状态。
- 事件用于可观察性,不作为唯一业务事实来源。
本章小结
Worker 事件是实例本地视角,QueueEvents 是跨 Worker 的队列视角;updateProgress() 用于可观察进度,处理函数返回值进入 returnvalue。关键结果必须在任务完成前可靠保存,不能把事件监听器当作事务提交点。