14 - 批量入队与吞吐优化

当一个请求需要创建几百或几千个任务时,逐个 await queue.add() 会产生大量 Redis 网络往返。BullMQ 提供批量 API,在减少往返次数的同时提供清晰的原子性边界。

1. 单队列批量添加:addBulk

import { Queue } from 'bullmq';

interface ThumbnailJobData {
  imageId: string;
  width: number;
}

const queue = new Queue<ThumbnailJobData>('thumbnails', {
  connection: { host: '127.0.0.1', port: 6379 },
});

const jobs = await queue.addBulk([
  {
    name: 'resize',
    data: { imageId: 'image-101', width: 320 },
  },
  {
    name: 'resize',
    data: { imageId: 'image-102', width: 640 },
  },
]);

console.log(`成功添加 ${jobs.length} 个任务`);

官方文档强调,单次 addBulk 对这批任务是“全部成功或全部失败”,并且通常比逐个添加更快。这里的原子性只覆盖这一次 BullMQ 入队操作,不会自动包含你的数据库事务或外部 API 调用。

2. 跨队列批量添加:FlowProducer.addBulk

Queue.addBulk 只能写入一个队列。需要原子地写入不同队列时,可以使用 FlowProducer.addBulk,即使这些任务没有父子依赖也可以。

import { FlowProducer } from 'bullmq';

interface EmailJobData {
  userId: string;
}

interface AuditJobData {
  action: string;
  userId: string;
}

const flowProducer = new FlowProducer({
  connection: { host: '127.0.0.1', port: 6379 },
});

const emailData: EmailJobData = { userId: 'user-42' };
const auditData: AuditJobData = {
  userId: 'user-42',
  action: 'registered',
};

await flowProducer.addBulk([
  {
    name: 'send-welcome-email',
    queueName: 'emails',
    data: emailData,
  },
  {
    name: 'write-audit-log',
    queueName: 'audits',
    data: auditData,
  },
]);

如果任务存在父子依赖,也可以把每棵 flow tree 放进同一次 flowProducer.addBulk()。整次调用同样是全部成功或全部失败。

3. 流水线与吞吐的正确理解

BullMQ 的设计目标之一是使用 Lua 脚本和 pipelining 提高 Redis 吞吐。对应用代码来说,最重要的结论不是手写 Redis pipeline,而是选择 BullMQ 已提供的高层批量 API:

不要绕过 BullMQ,直接对它的内部 Redis key 进行 pipeline 或 transaction。BullMQ 的 key 结构、Lua 脚本与状态转换共同维持队列一致性,手工修改可能破坏任务状态。

4. 为什么不能无限扩大单批大小

批量越大不一定越快。超大批次会带来:

更稳妥的做法是按固定上限切块,再逐批提交:

import { Queue } from 'bullmq';

interface IndexJobData {
  documentId: string;
}

const queue = new Queue<IndexJobData>('search-index', {
  connection: { host: '127.0.0.1', port: 6379 },
});

const BATCH_SIZE = 500;

async function enqueueDocuments(documentIds: string[]): Promise<void> {
  for (let start = 0; start < documentIds.length; start += BATCH_SIZE) {
    const batch = documentIds.slice(start, start + BATCH_SIZE);

    await queue.addBulk(
      batch.map((documentId: string) => ({
        name: 'index-document',
        data: { documentId },
        opts: {
          jobId: `index-${documentId}`,
          removeOnComplete: 1_000,
          removeOnFail: 5_000,
        },
      })),
    );
  }
}

500 只是便于理解的起点,不是官方规定。应根据任务数据大小、Redis 延迟、内存峰值和可接受的重试范围压测。

5. 批次的业务边界

技术批次和业务事务不是一回事。例如“导入 10 万个用户”可以切成 200 个技术批次,但业务上仍属于同一次导入。建议在任务数据中带上批次标识:

interface ImportJobData {
  importId: string;
  rowEnd: number;
  rowStart: number;
}

这样可以统计整个导入的进度,也能安全地重试某个技术批次。

6. 常见错误

逐个串行添加

for (const item of items) {
  await queue.add('process', item);
}

它最容易理解,但大量任务时网络往返明显。数量较多且需要同一原子边界时应改为 addBulk

无限制 Promise.all

await Promise.all(items.map((item) => queue.add('process', item)));

这会同时创建大量 Promise 和 Redis 请求,也不具备整批“全有或全无”的语义。小批并发可能有用,但不能替代 addBulk

把大文件直接放进 job.data

任务数据保存在 Redis 中。大二进制内容会增加网络、内存和持久化压力。更好的方式是把文件放入对象存储,只在任务中保存文件 ID、位置和校验信息。

7. 最佳实践清单

官方资料