batching

11. 批处理:攒一批,摊掉固定成本

0. 本章先解决什么问题

单次操作往往有一笔和数据量无关的固定成本:一次系统调用的进出内核、一次网络往返、一次 fsync。操作一多,固定成本累加起来比正事还贵。批处理的思路一句话:

攒一批一起做,固定成本除以批量数。
一次往返捎一百件事,每件事只摊 1% 的往返费。

先把各处的固定成本标上价,才知道每个批处理机制在省什么:

  • 一次系统调用:约 1 微秒(进出内核 + 上下文保存)
  • 一次同机房网络往返 RTT:约 0.5 毫秒
  • 一次 fsync:毫秒级(SSD 亚毫秒,机械盘 5~10 ms)
  • 一条 TCP/IP 报文的头部:固定几十字节,1 字节的数据也要背它
  • 一次网卡中断:微秒级 CPU 打断

本章按层扫过去:Redis 的 pipeline 一族、MySQL 的组提交和 JDBC 批量、Kafka 的攒批、TCP 的 Nagle、应用代码里的 buffer 和 StringBuilder,最后收拢批处理统一的反噬。

批处理的通用结构

这张图怎么读

上半部分对比逐个做和攒批做的成本结构:固定成本从每件一笔变成每批一笔。下半部分是各层的实例和它们各自”攒”的东西:命令、事务、消息、小包、字符。

1. Redis:pipeline、mget、事务、Lua 的精确区别

这四个东西都能”一次发多个命令”,但语义完全不同,放一起对比才记得牢:

往返账先算清: 单条命令的耗时构成里,命令执行本身微秒级,
网络往返 0.5 毫秒起,占比 90% 以上。
一万条命令逐条发 = 一万次 RTT ≈ 5 秒以上;
pipeline 打包发 = 几次 RTT + 一万次执行 ≈ 几十毫秒。
快一百倍的来源全在省往返,不在命令本身

省往返原子性中途读结果备注
pipeline无(命令间可插入别人的)不能,最后一起收纯客户端行为,攒着发
mget/mset单命令天然原子一次就是结果只限同类操作
事务 MULTI/EXEC排他执行不被插队不能,EXEC 才执行无回滚,入队错才废弃
Lua 脚本排他执行能,脚本里有逻辑能做”读了再决定”的复合操作

选型口诀:只是快就 pipeline;同类批量用 mget;要”读了再决定写什么”(比如库存判断再扣减)必须 Lua,pipeline 做不了,因为它发出去时还不知道前面命令的结果。
Lua 的板子在 Lua脚本,批处理与集群的注意事项在 zex - kv设计 & 批处理 & 主从or集群 & 配置(集群下 mget 的 key 要同槽,pipeline 要按节点分组,上上章的槽位机制在这里收费)。

一个纪律:pipeline 的批量别无限大。十万条命令一个批,服务端要攒十万条回复在输出缓冲里,客户端也要等全批结束才见第一条结果,通常几百到一千条一批。

2. MySQL:组提交和 JDBC 批量

2.1 组提交:fsync 的拼车

第 4 章预告过,这里展开。高并发下每个事务提交都单独 fsync,磁盘就是全局串行点。组提交把同一时刻在等的事务拼一车:

提交被拆成三段流水线(binlog 侧):
flush 段: 把这一批事务的 binlog 都写进文件(page cache)
sync 段: 一次 fsync 送整批落盘 <- 拼车发生在这
commit 段: 按序完成各自的引擎层提交
每段有队长(第一个到的事务),后来的搭车。
binlog_group_commit_sync_delay 可以故意等几微秒多凑几个
(和 Kafka 的 linger.ms 一个思路,见下节)
副产品: 同批提交的事务互不冲突这个事实被记进 binlog,
第 5 章的并行复制 LOGICAL_CLOCK 靠它判断从库能否并发重放

2.2 JDBC 批量:一个参数决定真批还是假批

代码块JAVA · 5 行收起展开
for (Order o : orders) {
    ps.setLong(1, o.getId());
    ps.addBatch();          // 攒
}
ps.executeBatch();          // 发

经典的坑:MySQL 驱动默认下,executeBatch 只是把语句挨个发,网络往返一次没省。连接串要加 rewriteBatchedStatements=true,驱动才把一批 INSERT 改写成一条多值语句:

INSERT INTO t VALUES (1,…),(2,…),(3,…)…
效果: 一次往返 + 一次解析 + 一次提交摊给全批,
万行插入从分钟级降到秒级
配套: 手动关自动提交,整批一个事务,
否则每行一个事务,fsync 次数没降下来

3. Kafka:把攒批做成产品的核心参数

Kafka 生产端的高吞吐一半来自攒批,两个参数控制”凑多少、等多久”:

batch.size(默认 16 KB): 同一分区的消息攒到这个量就发车
linger.ms(默认 0): 不满一批最多等多久。
默认 0 是延迟优先;吞吐型场景常设 5~20 ms,
故意等一等,换来批变大、请求数骤减
发车条件: 批满 或 时间到,先到先发(保险丝内置在参数里)
配套增益:
压缩按批做(compression.type): 一批里消息相似度高,
压缩率远好于逐条压,网络和磁盘一起省
上上章的粘性分区: 无 key 消息盯着一个分区攒,
攒满一批再换下一个,就是专门为攒批服务的路由策略
broker 端顺势收益: 整批写入、整批复制,
第 4 章的顺序写 + 页缓存吃到的都是成批的大块

4. TCP 层:Nagle、延迟确认和一个 40 毫秒的事故

内核也在替你攒批,不知道它的存在就会被它坑。

Nagle 算法(发送侧攒批):
小包别急着发,攒到一个 MSS 或等到在途数据被确认再发
动机: telnet 时代每敲一个字符发一个 41 字节的包,浪费惊人
延迟确认(接收侧攒批):
收到数据别立刻 ACK,等一会儿(最多 40 ms 量级),
看能不能搭上回程数据的车
两个好心凑一起的经典死锁:
发送方: 还有小包,但 Nagle 说要等 ACK
接收方: 延迟确认说 ACK 再等等
-> 双方互等,白白挂 40 ms
表现: 请求响应式的小包应用(RPC、Redis 协议)
偶发固定几十毫秒的延迟台阶
处方: 交互式协议关 Nagle(TCP_NODELAY)。
Redis 客户端、大多数 RPC 框架默认就关了,
Netty 里 childOption(ChannelOption.TCP_NODELAY, true) 是标配

教训有普遍性:多层各自攒批时,上层要么接管攒批(应用自己 pipeline,然后关掉内核的 Nagle),要么放手全交给下层,两层都自作主张就会出互等。

5. 应用代码里的攒批

BufferedOutputStream / BufferedReader(默认 8 KB 缓冲):
不带 buffer 逐字节 read/write,每字节一次系统调用,
八千倍的固定成本差距,最便宜的性能优化
Netty 的 write 和 flush 分家:
write 只进出站缓冲(攒),flush 才触发系统调用(发车)
writeAndFlush 逐条用是初学者常见的吞吐杀手,
批量场景应 write 多条再统一 flush
StringBuilder: 循环里 s += x 每次生成新 String(上一章的对象洪流),
StringBuilder 在可变工作区里攒,最后 toString 一次定型
ES 的 bulk、MyBatis 的 foreach 批量插入、
日志框架的异步 appender(攒日志批量刷盘),同一族

6. 反噬:批处理的四条账单

  1. 延迟换吞吐是明码标价
    linger.ms 设 20,每条消息平均多背 10 ms 延迟。
    吞吐型链路(日志、埋点)随便设,交易链路要掂量
  2. 保险丝必须双条件
    “攒满 N 或等满 T”缺一不可。只按量攒,低峰期最后一批
    永远凑不满,卡死在缓冲里(自己写业务攒批最容易漏 T)
  3. 批大小有上限
    批越大内存占用越大、失败重试的放大越狠、
    下游单次要消化的冲击越大(Kafka 的 max.request.size、
    pipeline 的回复缓冲都是这个约束的体现)
  4. 失败语义变复杂
    一批里第 7 条失败,前 6 条算什么?
    JDBC batch 部分失败的行为随驱动而异;
    Kafka 整批重试可能把成功过的消息再发一遍
    -> 批处理天然把”至少一次”塞给你,
    下游要么幂等要么去重,这条线索通往第 13 章的收尾

7. 联系实际:排查清单

  • 现象:循环调 Redis,量不大却整体很慢

  • 方向:逐条往返,RTT 占 90%;改 pipeline 或 mget

  • 现象:JDBC batch 用了,插入还是慢

  • 方向:没开 rewriteBatchedStatements,假批;顺手查自动提交关没关

  • 现象:RPC 偶发固定 40 ms 台阶

  • 方向:Nagle 撞延迟确认,确认 TCP_NODELAY

  • 现象:Kafka 生产吞吐上不去,请求数巨大

  • 方向:linger.ms=0 且消息小,批没攒起来;观察平均批大小指标

  • 现象:自研攒批服务低峰期数据总是最后差一截不落库

  • 方向:只有量条件没有时间条件,保险丝缺半根

  • 现象:Netty 应用 CPU 低但吞吐差,syscall 频繁

  • 方向:writeAndFlush 逐条刷,改成批 write + 单次 flush

8. 学完本章你能解决什么问题

  1. 各层的固定成本分别是什么量级,批处理各摊掉哪笔?
  2. pipeline、mget、事务、Lua 的四格差异表怎么填,何时必须 Lua?
  3. 组提交的三段流水线怎么拼车,副产品怎么喂给并行复制?
  4. rewriteBatchedStatements 改写前后差在哪,为什么要配手动提交?
  5. batch.size 和 linger.ms 怎么配合,粘性分区为什么是它的配套?
  6. Nagle 和延迟确认怎么互等出 40 ms,处方和通用教训是什么?
  7. Netty 的 write/flush 分家怎么用才对?
  8. 攒批的保险丝为什么必须双条件?
  9. 批处理为什么天然带来”至少一次”问题?

核心一句话:批处理用延迟换吞吐,把固定成本除以批量数;保险丝要”量满或时到”双条件,批大小要设上限,失败语义要想清楚;多层各自攒批时要指定谁做主,别让两层好心互等。

延伸阅读