05 - Job:数据、ID 与去重
Job 是队列中的一个任务实例。它不仅包含业务数据,还记录 ID、状态、执行次数、进度、返回值和失败原因。
1. 为任务数据和结果建模
import { Job, Queue, Worker } from "bullmq";
interface ThumbnailJobData {
imageId: string;
sourceKey: string;
width: number;
height: number;
}
interface ThumbnailJobResult {
thumbnailKey: string;
}
type ThumbnailJobName = "create-thumbnail";
const queue = new Queue<
ThumbnailJobData,
ThumbnailJobResult,
ThumbnailJobName
>("thumbnail");
const worker = new Worker<
ThumbnailJobData,
ThumbnailJobResult,
ThumbnailJobName
>(
"thumbnail",
async (
job: Job<ThumbnailJobData, ThumbnailJobResult, ThumbnailJobName>,
): Promise<ThumbnailJobResult> => {
return {
thumbnailKey: `thumbnails/${job.data.imageId}.jpg`,
};
},
);
泛型让生产者和 Worker 共享同一份数据契约。若任务种类较多,可以为每种任务定义独立类型,避免使用 any。
2. Job 数据经过 JSON 序列化
官方文档说明,任务数据写入 Redis 前会经过 JSON.stringify,Worker 端再通过 JSON.parse 还原。因此应传递普通 JSON 数据。
适合的类型包括:
- 字符串、数字、布尔值与
null; - 上述类型组成的数组;
- 只包含上述值的普通对象。
需要注意:
undefined、函数不会按原样保留;BigInt不能直接被 JSON 序列化;Map、Set、循环引用不能按预期保存;- 类实例会变成普通对象,原型方法、getter 和 setter 会丢失;
Date通常变成 ISO 字符串,Worker 需要显式还原。
推荐写法:
interface BillingJobData {
orderId: string;
amountInCents: number;
requestedAt: string;
}
await queue.add("create-thumbnail", {
imageId: "img-42",
sourceKey: "uploads/img-42.jpg",
width: 320,
height: 180,
});
金额使用最小货币单位的整数,时间使用 ISO 字符串,Worker 中再进行业务校验和转换。
3. 不要把大对象放入 Job
任务数据存储在 Redis 中。把文件二进制、超大 JSON 或完整数据库记录放进任务,会增加:
- 网络传输量;
- Redis 内存使用;
- 序列化和反序列化成本;
- 数据过期和隐私治理难度。
推荐只传稳定标识与必要参数:
interface VideoJobData {
videoId: string;
inputObjectKey: string;
preset: "mobile" | "desktop";
}
Worker 根据 videoId 或对象存储 key 获取最新数据。若任务必须处理入队时的快照,则要明确保存快照版本,避免读取到业务数据的意外新版本。
4. 进度与返回值
const worker = new Worker<ThumbnailJobData, ThumbnailJobResult>(
"thumbnail",
async (job): Promise<ThumbnailJobResult> => {
await job.updateProgress(10);
// 执行下载……
await job.updateProgress({ stage: "resizing", percent: 60 });
// 执行上传……
await job.updateProgress(100);
return { thumbnailKey: `thumbnails/${job.data.imageId}.jpg` };
},
);
processor 的返回值会保存到 job.returnvalue,也会随 completed 事件提供。
对于关键业务结果,官方建议在 processor 内部完成可靠持久化,而不要只依赖 completed 事件。因为“任务已完成”和“事件监听器写数据库”是两个操作,后者失败时任务仍可能已经进入 completed 状态。
5. 自定义 jobId
默认 ID 是队列范围内递增的计数器。也可以使用业务标识:
await queue.add(
"create-thumbnail",
{
imageId: "img-42",
sourceKey: "uploads/img-42.jpg",
width: 320,
height: 180,
},
{
jobId: "thumbnail-img-42-320x180",
},
);
同一队列内,如果该 ID 对应的 Job 仍然存在,再次添加会被忽略。这可用于简单的重复入队保护。
边界条件:
- 唯一性只在同一个队列内生效;
- 自定义 ID 不能包含
:; - 自定义 ID 不能是只包含数字的字符串,可添加业务前缀;
- Job 一旦被删除,相同 ID 可以再次添加;
- 因此
jobId不能代替数据库层的永久幂等约束。
6. BullMQ 内置去重
jobId 表达“这是同一个任务实例”,deduplication 则表达“一段时间内或某个任务完成前,只接受相同业务动作一次”。
Simple 模式
任务完成或最终失败前,相同去重 ID 的新任务会被忽略:
await queue.add(
"create-thumbnail",
{
imageId: "img-42",
sourceKey: "uploads/img-42.jpg",
width: 320,
height: 180,
},
{
deduplication: {
id: "thumbnail-img-42",
},
},
);
适合文件处理、同步操作等“同一资源在处理期间不要并发重复”的场景。
Throttle 模式
在固定 TTL 内忽略重复任务:
await queue.add("create-thumbnail", data, {
deduplication: {
id: `thumbnail-${data.imageId}`,
ttl: 5_000,
},
});
适合限制短时间内的重复触发。
Debounce 模式
在延迟窗口内不断用最新数据替换旧任务:
await queue.add("create-thumbnail", data, {
delay: 5_000,
deduplication: {
id: `thumbnail-${data.imageId}`,
ttl: 5_000,
extend: true,
replace: true,
},
});
适合搜索索引刷新、配置同步等“连续变化后只处理最后一次”的场景。
7. 去重不等于业务幂等
BullMQ 尽力提供一次语义,但最坏情况下仍可能至少执行一次,例如 Worker 在外部操作成功后、确认任务完成前崩溃。
涉及扣款、发券、创建订单等操作时,应在业务存储中建立幂等保护:
interface PaymentJobData {
paymentRequestId: string;
orderId: string;
amountInCents: number;
}
async function processPayment(job: Job<PaymentJobData>): Promise<void> {
// 数据库中应对 paymentRequestId 建唯一约束,
// 并在事务内判断该请求是否已经成功处理。
await chargeOnce(job.data);
}
Queue 侧去重减少重复任务,业务数据库的唯一约束防止重复副作用,两层保护解决的问题不同。
8. 更新 Job 数据要谨慎
官方提供 job.updateData() 更新已插入任务的数据,但初学阶段不要把 Job 当作频繁更新的数据库记录。更新时要考虑:
- Worker 是否已经读取旧数据;
- 重试是否应使用新数据;
- 审计时能否还原任务最初输入。
若需求是“保留最后一次变化”,优先评估 Debounce 去重模式,而不是多个进程随意修改 Job。
实用技巧与最佳实践
- 任务数据使用小型、扁平、可 JSON 序列化的 DTO。
- 在 Worker 边界再次校验数据,不要只依赖 TypeScript 编译期类型。
jobId使用稳定业务前缀,例如invoice-order_42,不要包含:。- 设计去重 ID 时,只选择真正决定“是否属于同一次工作”的字段。
- 关键结果在 processor 内可靠持久化;completed 事件更适合通知和观测。
- 数据中不要放密码、访问令牌或不必要的个人信息,因为任务可能被保留用于排错。
- 同时配置历史任务清理和数据库幂等,避免二者互相误解。
本章小结
Job 数据是跨进程、经过 JSON 序列化的消息契约。自定义 jobId 能避免仍存在任务的重复添加,deduplication 能表达时间窗口或执行期间的去重,但真正影响外部系统的操作仍必须由业务层保证幂等。