批处理与并发控制
队列里的消息不会一条一条送进消费者,那样握手和调度开销太大。Cloudflare Queues 会把消息攒成一批再投递,并且允许多个批次并发处理。怎么攒、攒多少、几条线并发,这些参数直接决定了消费速度和下游压力。这篇把批大小、批窗口、最大并发、消费限流四个参数讲透。
批处理的三个参数
消费者绑定时可以配置三个参数控制批次
[[queues.consumers]]
queue = "demo-queue"
max_batch_size = 10
max_batch_timeout = 5
max_concurrency = 2
各自含义
| 参数 | 作用 | 默认值 |
|---|---|---|
| max_batch_size | 一批最多多少条消息 | 10 |
| max_batch_timeout | 凑不满一批时最多等几秒就强制投递 | 0 |
| max_concurrency | 同时能跑几个批次 | 1 |
三个参数的协作逻辑。队列持续攒消息,攒到 max_batch_size 立即投递一个批次。如果消息来得慢,一直凑不满,就等 max_batch_timeout 秒后强制投递当前已有的。max_concurrency 控制同时能投出几个批次,并发越高处理越快但下游压力越大。
批大小怎么选
批大小影响单次处理成本和延迟
| 批大小 | 延迟 | 单条成本 | 适合场景 |
|---|---|---|---|
| 1 | 最低 | 最高 | 要求立即响应,比如推送通知 |
| 10 | 较低 | 较低 | 通用场景,均衡 |
| 100 | 较高 | 最低 | 写日志、批量入库等可延迟任务 |
批越大,每条消息分摊的调度开销越低,吞吐越高,但消息在队列里等凑批的时间也越长。对延迟敏感的任务用小批,对吞吐敏感的任务用大批。
代码里可以一次性处理整批,利用批量 API 降低下游调用次数
export interface Env {
DB: D1Database;
}
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
const stmts = batch.messages.map(msg => {
const body = msg.body as { uid: string; event: string };
return env.DB.prepare("INSERT INTO events (uid, event, ts) VALUES (?, ?, ?)")
.bind(body.uid, body.event, Date.now());
});
// 一次批量写入 D1(Cloudflare 的托管 SQLite 数据库),而不是逐条 INSERT
await env.DB.batch(stmts);
// 批量成功后统一 ack
for (const msg of batch.messages) msg.ack();
},
};
D1 的 batch 方法一次提交多条预编译语句,把整批消息一次性落库,比逐条 INSERT 快得多。
批窗口的作用
max_batch_timeout 是为低流量场景准备的。流量大时批次凑得快,这个参数几乎不触发。流量小时消息稀稀拉拉,没它的话最后几条消息会一直卡在队列里等凑满。
举个例子。max_batch_size 是 10,某分钟只来了 3 条消息。如果 max_batch_timeout 是 0,这 3 条会一直等到凑满 10 条才投递,可能等几分钟。设成 5 秒,那么第 5 秒就会强制投递这 3 条,不会无限等下去。
| 流量特征 | max_batch_timeout 建议 | 原因 |
|---|---|---|
| 高流量稳定 | 0 或较大 | 批次本来就凑得快,窗口意义不大 |
| 低流量偶发 | 1 到 5 秒 | 避免少量消息长时间滞留 |
| 实时性要求高 | 1 秒 | 优先保证延迟,牺牲一点批量优势 |
设成 0 表示禁用窗口,必须凑满才投递,要谨慎用。
最大并发的取舍
max_concurrency 决定同时有几个批次在跑。默认 1 是串行,一批处理完才投下一批。调大后多个批次并行,吞吐线性提升,但下游要扛得住。
| 并发数 | 吞吐 | 下游压力 | 风险 |
|---|---|---|---|
| 1 | 低 | 最小 | 几乎无风险 |
| 5 | 中高 | 中等 | 下游连接池要够 |
| 20 | 高 | 大 | 数据库连接可能打满,需限流 |
并发不是越大越好。假设下游是 Postgres,max_connections 是 100,每个批次要用 2 条连接,并发 20 就要 40 条连接,再叠加其他业务就可能打爆。
并发还会影响消息顺序。并发 1 时同一队列内消息严格按投递顺序处理。并发大于 1 时,先投递的批次可能后处理完,顺序不再保证。对顺序敏感的任务要么用并发 1,要么把消息按 key 分到不同队列。
消费限流
有时候下游能扛,但业务上不希望消费太快,比如调第三方 API 有速率限制。Queues 没有直接的速率限制参数,但可以用并发和批大小组合模拟
| 目标限流 | 配置组合 | 实际效果 |
|---|---|---|
| 每秒最多 10 条 | batch_size=10, timeout=1, concurrency=1 | 每秒一个 10 条批次 |
| 每秒最多 100 条 | batch_size=100, timeout=1, concurrency=1 | 每秒一个 100 条批次 |
| 每秒最多 200 条 | batch_size=100, timeout=1, concurrency=2 | 两个并发,每秒共 200 条 |
这是理论上限,实际还受处理耗时影响。更精确的限流要在代码里自己控制,比如用令牌桶
let lastTime = 0;
const MIN_INTERVAL = 100; // 每条至少间隔 100ms
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const msg of batch.messages) {
const now = Date.now();
const wait = MIN_INTERVAL - (now - lastTime);
if (wait > 0) await new Promise(r => setTimeout(r, wait));
lastTime = Date.now();
// 处理消息
msg.ack();
}
},
};
注意 setTimeout 在 Worker 里有 CPU 时间计费,长时间 sleep 不划算。限流场景更适合调小批大小和并发,让 Queues 自己控节奏。
参数调优流程
实际调参建议按这个顺序
第一步,先用默认值 max_batch_size=10、max_concurrency=1 跑起来,观察消费速度和下游负载。
第二步,如果消费跟不上生产,先调大 max_batch_size 到 50 或 100,看吞吐是否提升。
第三步,单批调到极限还不够,再调大 max_concurrency,每次加 1 观察下游。
第四步,低流量时段消息滞留,加 max_batch_timeout 到 1 到 5 秒。
| 现象 | 调整方向 |
|---|---|
| 消息堆积消费不动 | 加大 batch_size 或 concurrency |
| 下游报连接数超限 | 减小 concurrency |
| 低峰期消息延迟高 | 加大 batch_timeout |
| 消息顺序错乱 | concurrency 改回 1 |
小结
批处理四参数里,max_batch_size 管每批多少,max_batch_timeout 管凑不满时等多久,max_concurrency 管同时跑几批,组合起来就是消费节奏的全部。调参没有银弹,先默认值跑通,再按现象逐步调。下一篇讲消息出错了怎么办,死信队列和重试策略。