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 数据。

适合的类型包括:

需要注意:

推荐写法:

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 或完整数据库记录放进任务,会增加:

推荐只传稳定标识与必要参数:

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 仍然存在,再次添加会被忽略。这可用于简单的重复入队保护。

边界条件:

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 当作频繁更新的数据库记录。更新时要考虑:

若需求是“保留最后一次变化”,优先评估 Debounce 去重模式,而不是多个进程随意修改 Job。

实用技巧与最佳实践

本章小结

Job 数据是跨进程、经过 JSON 序列化的消息契约。自定义 jobId 能避免仍存在任务的重复添加,deduplication 能表达时间窗口或执行期间的去重,但真正影响外部系统的操作仍必须由业务层保证幂等。

官方资料