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 批量:一个参数决定真批还是假批
代码块收起展开
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. 反噬:批处理的四条账单
- 延迟换吞吐是明码标价
linger.ms 设 20,每条消息平均多背 10 ms 延迟。
吞吐型链路(日志、埋点)随便设,交易链路要掂量 - 保险丝必须双条件
“攒满 N 或等满 T”缺一不可。只按量攒,低峰期最后一批
永远凑不满,卡死在缓冲里(自己写业务攒批最容易漏 T) - 批大小有上限
批越大内存占用越大、失败重试的放大越狠、
下游单次要消化的冲击越大(Kafka 的 max.request.size、
pipeline 的回复缓冲都是这个约束的体现) - 失败语义变复杂
一批里第 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. 学完本章你能解决什么问题
- 各层的固定成本分别是什么量级,批处理各摊掉哪笔?
- pipeline、mget、事务、Lua 的四格差异表怎么填,何时必须 Lua?
- 组提交的三段流水线怎么拼车,副产品怎么喂给并行复制?
- rewriteBatchedStatements 改写前后差在哪,为什么要配手动提交?
- batch.size 和 linger.ms 怎么配合,粘性分区为什么是它的配套?
- Nagle 和延迟确认怎么互等出 40 ms,处方和通用教训是什么?
- Netty 的 write/flush 分家怎么用才对?
- 攒批的保险丝为什么必须双条件?
- 批处理为什么天然带来”至少一次”问题?
核心一句话:批处理用延迟换吞吐,把固定成本除以批量数;保险丝要”量满或时到”双条件,批大小要设上限,失败语义要想清楚;多层各自攒批时要指定谁做主,别让两层好心互等。