六维教程

消息队列基础

Workers(Cloudflare 的边缘计算函数服务)处理请求时,有些操作不需要立刻完成,比如发邮件、写日志、更新统计。如果让请求等着这些慢操作,响应时间会被拖长。Cloudflare Queues 就是为这种场景准备的,它把任务打包成消息放进队列,由另一个 Worker 慢慢消费,请求本身可以立即返回。这篇讲清楚消息队列的基础模型,以及如何在 Workers 里创建和使用一个队列。

生产者消费者模型

消息队列的核心是生产者消费者模型。三个角色分工明确

角色 职责 在 Queues 里的体现
生产者 把任务包装成消息投递到队列 任意 Worker 调用 send 方法
队列 暂存消息,按顺序等待被取走 Queues 托管的队列资源
消费者 从队列取消息并执行业务逻辑 绑定了队列的 Worker 的 queue 处理函数

整个链路

请求 Worker (生产者) -> 队列 -> 消费 Worker (消费者) -> 执行任务

请求 Worker 把消息丢进队列就返回,不用等消费完成。消费者按自己的节奏拉取并处理。两端解耦后,请求响应快,慢任务也不会丢。

什么时候该用队列

场景 同步处理的问题 用队列的好处
发送欢迎邮件 SMTP 慢,请求卡住 立即返回,后台发送
写用户行为日志 每条日志一次写库,请求变慢 批量攒起来一起写
视频转码、缩略图生成 耗时几十秒,请求早超时 异步处理,完成后通知
秒杀订单排队 瞬时高并发压垮数据库 削峰填谷,按数据库能力消费

创建队列

用 wrangler(Cloudflare 的命令行工具)创建队列

npx wrangler queues create demo-queue

列出已有队列

npx wrangler queues list

队列名全局唯一,创建后不能改名,只能删除重建

npx wrangler queues delete demo-queue

队列本身只是个消息暂存区,不绑定到任何 Worker 之前它不会做事。

生产者绑定

生产者是投递消息的 Worker。在 wrangler.toml 里加 producer 绑定

name = "producer-worker"
main = "src/index.ts"
compatibility_date = "2024-09-01"

[[queues.producers]]
binding = "DEMO_QUEUE"
queue = "demo-queue"

binding 是代码里的变量名,queue 指向真实队列名。代码里调用 send 投递消息

export interface Env {
  DEMO_QUEUE: Queue<string>;
}

export default {
  async fetch(request: Request, env: Env): Promise<Response> {
    const body = await request.json() as { email: string };
    // 把邮件任务丢进队列,不等待发送完成
    await env.DEMO_QUEUE.send({
      to: body.email,
      subject: "欢迎注册",
      sentAt: Date.now(),
    });
    return Response.json({ ok: true, msg: "已排队,稍后发送" });
  },
};

send 接受任意可序列化的 JSON 值,字符串、对象、数组都行。一次投一条也可以批量投递

await env.DEMO_QUEUE.sendBatch([
  { to: "a@x.com", subject: "欢迎注册" },
  { to: "b@x.com", subject: "欢迎注册" },
]);

sendBatch 一次最多 100 条消息,适合批量任务入库。

消费者绑定

消费者是真正干活的 Worker。绑定方式不同,用 queues.consumers

name = "consumer-worker"
main = "src/index.ts"
compatibility_date = "2024-09-01"

[[queues.consumers]]
queue = "demo-queue"
max_batch_size = 10
max_batch_timeout = 5

消费者 Worker 要导出一个 queue 函数,不是 fetch

export interface Env {
  // 消费者里可以绑定其他资源,比如数据库、KV
}

export default {
  async queue(batch: MessageBatch, env: Env): Promise<void> {
    for (const msg of batch.messages) {
      const body = msg.body as { to: string; subject: string };
      // 这里执行真正的发邮件逻辑
      console.log(`发送邮件到 ${body.to} 主题 ${body.subject}`);
      // 标记这条消息处理成功,从队列移除
      msg.ack();
    }
  },
};

batch.messages 是一批消息数组,逐条处理。msg.ack 表示成功,队列会删除这条消息。出错时用 msg.retry 让队列重新投递。

生产者和消费者可以是同一个 Worker,也可以是两个不同的 Worker。推荐拆开,职责清晰,单独扩缩容。

生产者消费者对照

两种绑定放一起对比更清楚

对比项 producers 绑定 consumers 绑定
代码入口 fetch 函数 queue 函数
作用 投递消息 处理消息
触发方式 收到 HTTP 请求时调用 send 队列有消息时自动拉取
必填字段 binding、queue queue
可选字段 max_batch_size、max_batch_timeout 等

一个队列可以有多个生产者投递,但只能有一个消费者 Worker 绑定。如果要多个 Worker 处理同一队列,目前做不到,得在消费端自己分发。

本地开发调试

wrangler 支持本地模拟队列。启动 dev 时生产者的 send 不会真的投递到线上,而是先进本地缓冲,再触发消费者的 queue 函数

npx wrangler dev

观察消息流转用 wrangler tail 看实时日志

npx wrangler tail consumer-worker

也可以手动往队列塞消息做测试

npx wrangler queues send demo-queue '{"to":"test@x.com","subject":"测试"}'

这条命令直接往队列写一条消息,会触发消费者处理,适合不启动生产者就能验证消费逻辑。

小结

消息队列把请求和慢任务解耦,生产者投递消息立即返回,消费者按节奏处理。在 Workers 里就是两个绑定,producers 让 Worker 能 send,consumers 让 Worker 自动消费。下一篇讲批处理与并发控制,把消费节奏调到最优。

下一篇 批处理与并发控制

上一篇
Pages Functions 中间件
下一篇
批处理与并发控制