六维教程

批处理与并发控制

队列里的消息不会一条一条送进消费者,那样握手和调度开销太大。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 管同时跑几批,组合起来就是消费节奏的全部。调参没有银弹,先默认值跑通,再按现象逐步调。下一篇讲消息出错了怎么办,死信队列和重试策略。

上一篇 消息队列基础
下一篇 死信队列与重试策略

上一篇
消息队列基础
下一篇
死信队列与重试策略