《TypeScript编程实战》9.1 BullMQ 队列模型与 payload 泛型

本节从「为什么要把工作丢进队列」讲起,拆解 BullMQ 的 Queue、Worker、Job 三个核心对象及其在 Redis 中的存储结构,再重点解决一个工程痛点:如何用泛型把 payload 类型从生产者贯穿到消费者。你会学到 Queue 的类型参数推导、job name 的字面量联合、判别联合注册表的写法,以及 Date 序列化等常见坑。读完能把队列调用从 any 变成端到端可检查。

本节目标:搞清楚任务队列在真实系统里解决的是哪一类问题,掌握 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 重试、幂等与死信 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「typescript」更多文章

  1. 《TypeScript高级编程》11.3 类型驱动架构与团队规范
  2. 《TypeScript高级编程》11.2 渐进式迁移与严格化路径
  3. 《TypeScript高级编程》11.1 TS 版本演进与 breaking changes