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:
- 多个同队列任务:
queue.addBulk()。 - 多队列任务或多棵依赖树:
flowProducer.addBulk()。 - 单个普通任务:
queue.add()。
不要绕过 BullMQ,直接对它的内部 Redis key 进行 pipeline 或 transaction。BullMQ 的 key 结构、Lua 脚本与状态转换共同维持队列一致性,手工修改可能破坏任务状态。
4. 为什么不能无限扩大单批大小
批量越大不一定越快。超大批次会带来:
- 更大的请求序列化与网络包。
- Redis 执行单次操作时间变长。
- Node.js 内存峰值升高。
- 一次失败需要重试更大的范围。
- 更难定位哪条源数据不合法。
更稳妥的做法是按固定上限切块,再逐批提交:
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. 最佳实践清单
- 批量前先验证输入,避免一条坏数据使整批失败。
- 给可重复提交的任务设置稳定
jobId或其他幂等键。 - 批次大小通过压测决定,并记录每批耗时与失败率。
- 为 completed/failed 任务配置合理的自动清理策略。
- 生产者要有背压:队列积压超过阈值时暂停继续灌入。
- 跨 Redis Cluster 队列执行原子 bulk/flow 时,先确认相关 key 位于兼容的哈希槽,详见生产章节。
- 关闭应用时调用
queue.close()和flowProducer.close()。