消息队列基础
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 自动消费。下一篇讲批处理与并发控制,把消费节奏调到最优。
下一篇 批处理与并发控制