本节目标:搞清楚任务队列在真实系统里解决的是哪一类问题,掌握 BullMQ 的 Queue / Worker / Job 模型,并且把「payload 类型只写在注释里」这种常见做法换成由编译器强制校验的泛型方案。读完本节,你应该能独立搭起一个生产者与消费者共享同一套 payload 类型的后台任务链路。
9.1 BullMQ 队列模型与 payload 泛型
先看一个几乎每个后端都写过的注册接口:
// 反面教材:把重活留在请求链路里
app.post('/register', async (req, res) => {
const user = await prisma.user.create({ data: req.body });
await sendWelcomeEmail(user); // 200ms
await generateAvatar(user); // 800ms
await syncToCrm(user); // 1200ms
res.json({ id: user.id });
});
这段代码没有语法错误,类型也全是干净的。但用户会等 2 秒以上,而其中任何一个下游抖动都会让注册接口整体 500——邮件服务挂了,用户就注册不了。把这三件事挪进队列,是本节要讲的全部动机。
9.1.1 队列解决的三个问题
| 问题 | 同步调用的表现 | 队列之后的表现 |
|---|---|---|
| 削峰 | 秒杀流量直接打穿下游 | 请求只写一条 job,worker 按自身速率消费 |
| 解耦 | 邮件服务故障导致注册失败 | 邮件故障只影响邮件,注册照常成功 |
| 可靠重试 | 失败即丢,只能靠用户重试 | job 持久化在 Redis,失败自动重试 |
第三条最容易被低估。同步调用里,一次网络抖动就是一次用户可见的失败;队列里,它只是一次重试。而重试要成立,前提是 job 被持久化了——这正是 BullMQ 与内存队列(比如 p-queue)的根本差别。
9.1.2 Queue、Worker、Job 三个对象
BullMQ 的 API 表面很薄,就三个角色:
| 对象 | 职责 | 通常部署在哪 |
|---|---|---|
Queue | 生产端:add() 投递 job,查询统计 | Web 进程 |
Worker | 消费端:注册 processor,取出并执行 job | 独立的 worker 进程 |
QueueEvents | 事件订阅:监听 completed / failed 等 | 需要感知结果的一方 |
关键点是:生产者与消费者是两套进程。它们之间唯一的契约就是「job name + payload 结构」。这个契约如果只写在文档或注释里,改错字段名要到线上跑挂了才发现——本节剩下的篇幅都在解决这一件事。
在 Redis 里,BullMQ 用一组键来维护状态,理解它有助于排查问题:
bull:email:wait # list:等待被消费的 job id
bull:email:active # list:正在被消费的 job id
bull:email:delayed # zset:延迟 job,score 是执行时间戳
bull:email:completed # zset:已完成的 job id(受 removeOnComplete 控制)
bull:email:failed # zset:失败的 job id,死信就从这里来
bull:email:1 # hash:单个 job 的数据(name / data / opts)
bull:email:events # stream:事件流,QueueEvents 消费的就是它
9.1.3 最小可运行示例
pnpm add bullmq ioredis
// queue.ts —— 生产者
import { Queue } from 'bullmq';
import IORedis from 'ioredis';
const connection = new IORedis({ host: '127.0.0.1', port: 6379, maxRetriesPerRequest: null });
export const emailQueue = new Queue('email', { connection });
await emailQueue.add('welcome', { to: 'ada@example.com', name: 'Ada' });
// worker.ts —— 消费者(独立进程启动)
import { Worker } from 'bullmq';
import IORedis from 'ioredis';
const connection = new IORedis({ host: '127.0.0.1', port: 6379, maxRetriesPerRequest: null });
new Worker(
'email',
async (job) => {
console.log(job.name, job.data); // 这里 job.data 是 any
return { sent: true };
},
{ connection },
);
$ pnpm tsx worker.ts
welcome { to: 'ada@example.com', name: 'Ada' }
注意上面那行注释:job.data 是 any。to 拼错成 too 不会有任何提示,job.data.user.name 这种写法也不会被拦。这就是 BullMQ 默认类型参数埋下的坑——它把类型安全的选择权留给了调用方。
9.1.4 Queue 的三个类型参数
Queue 的签名(简化后)是这样的:
declare class Queue<
DataType = any,
ResultType = any,
NameType extends string = string,
> {
add(name: NameType, data: DataType, opts?: JobsOptions): Promise<Job<DataType, ResultType, NameType>>;
}
三个参数分别对应「payload」「processor 返回值」「job name」。默认全是宽类型,所以只要显式传第一个,整条链路的推导就活了起来:
export interface EmailPayload {
to: string;
name: string;
template: 'welcome' | 'reset-password' | 'invoice';
}
export interface EmailResult {
messageId: string;
}
export const emailQueue = new Queue<EmailPayload, EmailResult>('email', { connection });
await emailQueue.add('welcome', {
to: 'ada@example.com',
name: 'Ada',
template: 'welcome',
});
此时若漏字段或写错模板名,编译期就会报:
error TS2345: Argument of type '{ to: string; name: string; template: "welcom"; }'
is not assignable to parameter of type 'EmailPayload'.
Types of property 'template' are incompatible.
Type '"welcom"' is not assignable to type '"welcome" | "reset-password" | "invoice"'.
9.1.5 消费者侧的类型对齐
消费者必须手动声明同样的类型参数,因为 Worker 拿不到生产者的 Queue 实例:
import { Worker, type Job } from 'bullmq';
new Worker<EmailPayload, EmailResult>(
'email',
async (job: Job<EmailPayload, EmailResult>) => {
const { to, name, template } = job.data; // 全部有类型
await mailer.send({ to, template, vars: { name } });
return { messageId: crypto.randomUUID() };
},
{ connection, concurrency: 5 },
);
工程上的做法是把 payload 与 result 类型抽到一个共享包里(比如 monorepo 中的 packages/contracts),生产者和消费者都从那里 import。这样「两边类型漂移」就不再靠人盯——只要一方改了接口,另一方的 pnpm typecheck 立刻红。
9.1.6 一个队列放多种任务:name 的字面量联合
第三个类型参数 NameType 常被忽略,但它是「同一队列多任务」的关键:
export const emailQueue = new Queue<EmailPayload, EmailResult, 'welcome' | 'reset-password'>('email', {
connection,
});
await emailQueue.add('welcome', payload); // OK
await emailQueue.add('reset', payload); // 报错:'"reset"' 不在联合里
看起来很美好,但有个硬伤:NameType 只约束了 name,没有把 name 和 data 关联起来。你依然可以 add('welcome', 重置密码的 payload)。要让 name 与 payload 一一对应,需要换一种建模方式。
9.1.7 判别联合 payload 注册表
思路是:把「job name → payload 类型」写成一个映射表,再让 name 变成 payload 的判别属性。
// contracts/jobs.ts
export interface JobMap {
'email:welcome': { to: string; name: string };
'email:reset': { to: string; token: string; expiresAt: string };
'report:generate': { userId: string; from: string; to: string };
}
export type JobName = keyof JobMap;
然后写一个薄薄的泛型包装,把 add 收窄:
// queue.ts
import { Queue, type Job } from 'bullmq';
const raw = new Queue('tasks', { connection });
export function enqueue<N extends JobName>(
name: N,
data: JobMap[N],
opts?: Parameters<typeof raw.add>[2],
): Promise<Job> {
return raw.add(name, data, opts);
}
enqueue('email:welcome', { to: 'ada@example.com', name: 'Ada' }); // OK
enqueue('email:welcome', { to: 'ada@example.com', token: 'x' }); // 报错:缺少 name
enqueue('email:reset', { to: 'ada@example.com', token: 'x', expiresAt: '2026-10-01' }); // OK
推导结果可以这样验证:
type A = Parameters<typeof enqueue<'email:reset'>>[1];
// ^? type A = { to: string; token: string; expiresAt: string }
消费侧则按 name 分发,用 switch 把 payload 窄化:
// worker.ts
type Handlers = { [N in JobName]: (data: JobMap[N]) => Promise<unknown> };
const handlers: Handlers = {
'email:welcome': async (d) => mailer.sendWelcome(d.to, d.name),
'email:reset': async (d) => mailer.sendReset(d.to, d.token),
'report:generate': async (d) => reporter.build(d.userId, d.from, d.to),
};
new Worker('tasks', async (job) => {
const name = job.name as JobName;
const handler = handlers[name] as (d: unknown) => Promise<unknown>;
return handler(job.data);
});
这里最后两行的断言是不可避免的妥协:BullMQ 无法知道 Redis 里那条 job 究竟是哪个 name,所以边界处必须有一次断言。把断言集中在这一个函数里,是「类型安全」与「运行时现实」之间的标准取舍——边界收敛,而不是散落各处。
9.1.8 五个常见坑
一、JSON 序列化会吃掉 Date。 BullMQ 用 JSON 存 payload,Date 会变成字符串:
await enqueue('email:reset', { to: 'a@b.com', token: 't', expiresAt: new Date() });
// 消费者拿到的 expiresAt 实际是 "2026-09-27T02:00:00.000Z",但类型仍写着 Date
这是类型撒谎的典型场景:JobMap 里写 Date 编译不报错,运行时 d.expiresAt.getTime() 直接抛 TypeError: d.expiresAt.getTime is not a function。规范做法是契约里一律用 ISO 字符串,或在消费端用 zod 反序列化。
二、payload 不宜过大。 一个 job 的 data 存在 Redis hash 里,建议控制在 100KB 以内。大对象(比如整份报表数据)应该只传 id,让 worker 自己去数据库取。
三、循环引用会直接抛错。 JSON.stringify 遇到循环引用抛 TypeError: Converting circular structure to JSON,而 Prisma 的某些关联对象很容易带环。
四、不要用 as Job<MyPayload> 兜底。 这类断言把错误从编译期推迟到运行期,等价于放弃类型。正确做法是在 enqueue 这一层收窄。
五、队列名前缀与多环境。 本地调试别连生产 Redis;用 prefix: 'bull:dev' 或独立实例隔离,否则本地 worker 会真的把生产的邮件发出去。
9.1.9 与本书其它章节的衔接
payload 类型的来源往往是数据层:Prisma 生成的类型可以直接作为契约基础,参见 《TypeScript编程实战》7.1 Prisma schema 与类型生成 。连接与封装的规范写法见 《TypeScript编程实战》8.2 Redis 类型安全封装 。而 job 执行过程的可观测性(结构化日志、追踪 id 透传)见 《TypeScript编程实战》3.3 结构化日志与脱敏 。
站内已有专题对 BullMQ 做过单点深挖,可作为延伸阅读:TypeScript 任务队列与 BullMQ 、Node.js BullMQ 后台任务 。若想横向比较不同队列的取舍,可读 消息队列方案对比 。
小结
本节把「重活丢进队列」这件事拆成了三个层次。模型层:BullMQ 只有 Queue、Worker、Job 三个角色,生产者与消费者是两套进程,它们之间唯一的契约就是 job name 与 payload 结构。类型层:Queue 的三个类型参数默认都是 any,必须显式传入才有推导;NameType 能约束名字但无法把名字与 payload 绑定,要真正做到一一对应得靠判别联合注册表。边界层:JSON 序列化会吃掉 Date,所以契约里应当用字符串;Redis 与 BullMQ 之间的那一次类型断言无法避免,但必须收敛在一个函数里。
最容易犯的错是把类型写在注释里——它看起来像文档,实际上没有任何强制力。把契约抽成共享包、让生产者与消费者都 import 同一份类型,是这个问题的正解。
投递只是开始:job 失败之后怎么办,是本节的下一半。接下来 《TypeScript编程实战》9.2 重试、幂等与死信 会讲清退避策略、幂等键的设计,以及如何把反复失败的 job 送进死信队列而不是无限重试。
阅读导航:上一节:8.3 穿透·击穿·雪崩防护 · 下一节:9.2 重试、幂等与死信 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。