04 - Queue 与生产者
Queue 代表一条任务队列。生产者通过 Queue.add() 把业务意图转换为任务,Redis 保存任务,Worker 可以现在或稍后处理它。
1. 创建类型明确的 Queue
import { Queue } from "bullmq";
interface ReportJobData {
userId: number;
format: "csv" | "pdf";
}
interface ReportJobResult {
fileKey: string;
}
type ReportJobName = "generate-report";
const reportQueue = new Queue<
ReportJobData,
ReportJobResult,
ReportJobName
>("report", {
connection: {
host: "127.0.0.1",
port: 6379,
},
});
三个泛型依次描述任务数据、Worker 返回值与任务名称。这样可以在编译阶段发现字段或任务名拼错。
创建 Queue 时,BullMQ 会创建或复用相应的队列元数据。创建同名 Queue 不会清空已有任务。
2. 使用 add 添加任务
const job = await reportQueue.add("generate-report", {
userId: 42,
format: "pdf",
});
console.log(job.id);
add 的三个主要参数是:
- 任务名称;
- 可 JSON 序列化的数据;
- 可选的任务选项。
任务名称用于区分同一队列中的业务动作。对于处理资源和失败策略完全不同的任务,通常更适合拆成不同队列。
3. 常用任务选项
await reportQueue.add(
"generate-report",
{ userId: 42, format: "pdf" },
{
attempts: 3,
backoff: {
type: "exponential",
delay: 1_000,
},
delay: 5_000,
removeOnComplete: 100,
removeOnFail: 1_000,
},
);
attempts:最大执行次数,包含第一次执行;backoff:失败重试前的等待策略;delay:至少等待指定毫秒数后才有资格执行;removeOnComplete、removeOnFail:限制保留的历史任务数量。
延迟不是精确计时器。到达延迟时间只代表任务可以被处理,实际开始时间还受 Worker 是否空闲、并发限制等影响。
4. 使用 defaultJobOptions
把队列级通用策略集中设置:
const reportQueue = new Queue<ReportJobData>("report", {
connection: {
host: "127.0.0.1",
port: 6379,
},
defaultJobOptions: {
attempts: 3,
backoff: {
type: "exponential",
delay: 1_000,
},
removeOnComplete: {
age: 3_600,
count: 1_000,
},
removeOnFail: {
age: 24 * 3_600,
count: 5_000,
},
},
});
单次 add 提供的选项可以覆盖默认值。适合放在默认选项中的内容是全队列一致的保留和重试策略;与具体业务相关的 jobId、延迟或优先级更适合在添加任务时指定。
5. 批量添加任务
const jobs = await reportQueue.addBulk([
{
name: "generate-report",
data: { userId: 1, format: "csv" },
},
{
name: "generate-report",
data: { userId: 2, format: "pdf" },
},
]);
console.log(`成功添加 ${jobs.length} 个任务`);
官方文档指出,addBulk 对同一队列具有“全部成功或全部失败”的原子性,并通过减少 Redis 往返次数提高批量写入性能。
如果需要跨多个队列原子添加,可学习后续 FlowProducer 的 addBulk,不要把多次普通 add 误认为跨队列事务。
6. 暂停与恢复队列
await reportQueue.pause();
// 执行维护操作……
await reportQueue.resume();
Queue 的 pause() 是全局暂停:所有 Worker 都不再领取新任务,但已经开始的任务会继续执行到成功或失败。任务仍然可以被加入暂停中的队列。
如果只想停止某一个 Worker,应调用 worker.pause(),不要暂停整个队列。
7. 清理任务的三个层级
自动保留策略
生产环境优先使用 removeOnComplete 和 removeOnFail 控制历史数据。官方建议通常少保留一些成功任务,多保留一些失败任务,便于排错。
drain
await reportQueue.drain(true);
清除等待任务;传入 true 时也清除延迟任务。它不会清除 active、completed、failed 等状态的任务。
clean
const deletedJobIds = await reportQueue.clean(
60_000,
1_000,
"completed",
);
删除某个状态中早于宽限期的任务,并限制单次删除数量。分批清理比一次处理大量历史任务更容易控制 Redis 压力。
obliterate
obliterate() 会删除整个队列及其内容,属于破坏性维护操作。不要把它写进普通业务流程;执行前应确认队列名、暂停生产者,并做好环境隔离。
8. 把 Queue 当作长期对象
错误做法是在每个 HTTP 请求中创建 Queue:
// 不推荐:高频请求会反复创建对象和连接。
async function handleRequest(): Promise<void> {
const queue = new Queue("report");
await queue.add("generate-report", { userId: 42, format: "pdf" });
await queue.close();
}
更好的做法是在应用启动时创建一次,在路由中复用,并在进程退出时关闭。
实用技巧与最佳实践
- 把队列名和任务名集中定义,避免生产者与 Worker 不一致。
- 为 Queue 添加 TypeScript 泛型,不要让任务数据退化为不明确的对象。
- 对瞬时网络故障使用有限重试和退避;对参数错误等永久失败不要盲目重试。
- 为完成和失败任务设置保留上限,否则 Redis 会持续增长。
- 批量导入优先使用
addBulk,不要在循环中串行执行大量add。 - 清理操作应放在受控的运维脚本中,并记录删除数量。
- 入队成功只代表任务已保存,不代表业务处理成功。
本章小结
生产者的核心职责是正确地构造任务并可靠入队。defaultJobOptions 统一策略,addBulk 提升批量写入效率,暂停和清理 API 用于运维;它们都不应掩盖业务层对失败、幂等和历史保留的设计。