阶段二 · 架构原理与调优

生产调优:批量,压缩,linger.ms 与 buffer.memory

前置回顾

  • 第 03 篇给过 Producer 关键参数速查表,当时是"认识每个旋钮",本篇讲"旋钮之间怎么联动"
  • 第 06 篇的可靠性配置(acks=all + 幂等)与第 09 篇的副本确认,决定了每次请求要等的延迟下限
  • 第 08 篇讲过压缩 batch 让 Broker 原样转发、保住零拷贝,本篇从发送端看压缩的收益与成本

一句话总结

Producer 调优只有一个核心交换:批量换吞吐,延迟是代价.
batch.size 和 linger.ms 决定批量怎么形成,压缩把批量收益再放大一遍,buffer.memory 是兜底的背压闸门.
调参的第一步不是改参数,是定目标:这套业务到底是延迟优先还是吞吐优先.

先定目标:吞吐与延迟的取舍

一条消息从 Produce() 到回调被触发,延迟由四段组成:

Produce() 入队
│
├─ 1) 等凑批: 0 ~ linger.ms ← 本篇的主角, 唯一可策略性伸缩的一段
│
├─ 2) 网络往返: 一次请求的 RTT
│
├─ 3) Broker 处理: 追加写入 + (acks=all 时)等 ISR 确认(第 09 篇)
│
└─ 4) 返回 ack → callback
  1. 和 4) 是物理约束,3) 由可靠性配置决定(第 03/06 篇),只有 1) 是可以主动调大来换吞吐的.吞吐的构成:
吞吐 = 每请求携带的消息数 × 每秒请求数
└── 批量大小 ──┘ └─ 受 RTT 与在途数限制, 提升空间有限 ─┘

所以整篇的脉络很清晰:先把批量做大(唯一的大杠杆),再用压缩把批量的收益放大,最后用 buffer 和背压控制"批量做大"的副作用.

目标 典型业务 调参倾向
延迟优先 在线交易、实时风控 小批量、立即发送,宁可牺牲吞吐
吞吐优先 日志采集、数据管道 大批量、等一段时间再发
两者都要 大多数业务 高流量靠批量自然形成,低流量不额外等待

批量:吞吐的第一杠杆

批是怎么攒起来的

Producer 内部有个 RecordAccumulator,每个分区一条 batch 队列.业务代码调 Produce() 只是把消息追加进当前批次,真正的发送由 Sender 线程负责:

RecordAccumulator

Partition 0: [batch 已满, 待发] [batch 攒批中]
Partition 1: [batch 攒批中]

Sender 循环:
1. 挑出可发送的 batch
2. 按目标 Broker 聚组(同一 Broker 上多个分区的 batch 合并进一次请求)
3. 发出, 等 ack

哪些条件触发"这批发出去":

触发条件 说明
batch 满了 攒到 batch.size 上限,立即发
linger 到期 linger.ms 时间到,哪怕只攒了 3 条也发
buffer 耗尽 内存池没空间了,强制把最大的 batch 发出去腾位置
手动 Flush 调 Flush(),用于关停等场景

批量为什么能放大吞吐

一次网络请求的往返与处理开销是固定的,批量让这份固定开销摊到 N 条消息上:

每次请求的往返 + 处理按 0.5ms 粗估, 单连接串行请求上限约 2000 次/s

逐条发送: 2000 次/s × 1 条 = 约 2000 条/s
每批攒 500 条: 2000 次/s × 500 条 = 约 100 万条/s

同样的物理开销, 批量把吞吐放大了几百倍

批量发送 = 攒满一车再发车:单件快递随到随发,司机跑断腿;攒够一车发一次,同样一趟运输量翻几百倍.
linger.ms = 发车时刻表:哪怕车没装满,到点也发,这是延迟的兜底.


linger.ms 与 batch.size:谁先到谁触发

这两个参数常被混在一起说,职责其实正交:

  • batch.size:一批最多攒多大(上限)
  • linger.ms:一批最多等多久(超时)

谁先满足谁触发.组合起来的行为:

组合 低流量时 高流量时 适用
小 batch + linger 0 每条立即发 小批快发,吞吐上不去 极少推荐
大 batch + linger 0 每条立即发 自动攒满大批 通用默认:高流量吃批量红利,低流量不吃延迟
大 batch + linger 5-50ms 等 linger 到点凑批发 攒满即发 吞吐优先,接受固定的延迟底噪

"大 batch + linger 0" 的隐藏优势值得单独说:静止时消息立即走(不牺牲延迟),洪峰时 batch 自然被填满(吃满吞吐红利).流量越是波峰波谷分明,这个组合越划算.

而 linger.ms > 0 的真实代价是给每条消息叠加了 0 ~ linger 的等待,延迟预算紧张的服务要把它算进 P99 里:

P99 预算 100ms:
等凑批 20ms + 网络与副本确认 30ms + 业务余量 50ms ← linger 已经吃掉两成预算

linger 调大了,实际延迟可能远不止 linger

linger 只是"最多等多久",实际延迟还取决于 Sender 线程是否及时取走 batch.
发送压力大、Sender CPU 饱和时,到期的 batch 要排队等发送,端到端延迟会成倍于 linger 值.
调大 linger 前,先确认 Sender 不是瓶颈(客户端 CPU 与第 07 篇的 request-latency-avg 指标).


压缩:压缩单位是 batch

为什么批量越大,压缩越划算

压缩的对象是整个 batch,不是单条消息.消息之间重复的字段名、相似的结构,在整批范围内才有压缩空间;一条消息单独压,几乎压不动.

压缩 = 真空收纳袋:散装衣物直接装箱占满空间,抽真空后同样一箱多装几倍.
抽一次袋子的固定成本摊在整袋衣物上,所以整批一起压比一件件压划算得多.

压缩发生在 Producer 端,之后的链路全部受益:

Producer 压缩整批 ──► Broker 原样存储(不解压)──► Consumer 拉取后解压
│
├─ 两段网络带宽都省
├─ Broker 磁盘占用变小
└─ 第 08 篇的 sendfile 零拷贝照常工作(不解压才能直传)

算法怎么选

compression.type 压缩率 CPU 开销 说明
none 1x 无 Java 客户端默认;franz-go 也不默认压缩
lz4 中 很低 吞吐场景首选,生态最成熟
snappy 中 低 老牌算法,逐渐被 lz4 取代
zstd 高 中 2.1+ 支持,压缩率接近 gzip 但快得多,级别可调
gzip 高 高 兼容性最广,性能最差,新系统没必要选

具体压缩率取决于消息内容的重复度,文本类消息(JSON 日志)通常能到 2-4 倍,已经压缩过的二进制(图片、Protobuf 密文)收益有限.选择策略:

  • 拿不准就 lz4:CPU 便宜,带宽立省一半左右
  • 带宽或磁盘是明确瓶颈,上 zstd 并把级别从默认的 3 往上试
  • 压缩 CPU 花在客户端:高吞吐 Producer 要给 Sender 留出 CPU,容器 CPU limit 要按压缩后的成本重新评估

broker 端 compression.type 别乱动

Topic 级和 Broker 级也有 compression.type 配置,默认值 producer 表示"保持 Producer 压缩的格式原样存储",这是最优路径.
改成 none 或指定算法,会让 Broker 解压再重压缩:CPU 白烧,还破坏第 08 篇的零拷贝转发.压缩决策留在 Producer 端.


buffer.memory:背压与内存水位

它在防什么

buffer.memory 是 RecordAccumulator 的内存池上限,所有分区的 batch(攒批中的 + 已发送等 ack 的)共用这个池子.它的作用是解耦:

生产速率 > 网络消化速率 时:

业务代码 ──Produce()──► [ buffer 池 ] ──► Sender ──► Broker
│
└─ 池满 → Produce() 阻塞(或等待超时)
→ 把压力传回业务代码 ← 这就是背压

没有这个池子,要么每条消息阻塞在网络上(吞吐崩溃),要么无限内存增长(进程 OOM).有了池子,正常波动被吸收,极端情况把压力传导回上游.

池满时的行为:

客户端 行为
Java send() 阻塞,超过 max.block.ms(默认 60s)抛 TimeoutException
franz-go Produce() 阻塞,直到有空间,或 context 被取消

容量怎么估

buffer 最低容量 ≈ 峰值写入速率 × 可接受的阻塞时长

写入 50 MB/s, 允许最多阻塞 2s:
→ 至少 100MB

Java 默认 32MB,小水管够用,高吞吐场景远远不够.franz-go 的对应旋钮有两个维度:

选项 默认 作用
kgo.MaxBufferedRecords 10000 条 按条数限流
kgo.MaxBufferedBytes 按需配置 按字节限流

按条数还是按字节,取决于消息大小是否均匀:大小固定按条数就够,大小差异大要按字节.

buffer 满的阻塞会级联到业务

Produce() 阻塞时,占用的是业务 goroutine(或 Java 线程).HTTP handler 里同步发消息,就会表现为接口超时.
对策:发消息路径设超时与降级(第 05 篇的错误分级),或用异步 + 有界队列先把业务逻辑与发送解耦.


在途请求:把 RTT 藏起来

前面的批量解决"每次请求带多少条",还有一个正交的维度:同时有几个请求在途.

串行:  请求1 ──等 ack──► 请求2 ──等 ack──► 请求3      吞吐 ≈ 批量 / RTT
管道: 请求1 ──等 ack──►
请求2 ──等 ack──► 同时在途, 互不等待 吞吐 × 在途数

max.in.flight.requests.per.connection 控制这个数:

  • 非幂等模式:可以调很大,但重试会乱序(第 03 篇的坑)
  • 幂等模式:上限 5,由 Broker 按 sequence number 去重排序,顺序仍有保证、吞吐照样管道化

franz-go 对应 kgo.MaxProduceRequestsInflightPerBroker,幂等开启时默认 5.这解释了第 06 篇"acks=all 也能高吞吐"的组合拳:每批等 ISR 确认确实慢,但 5 个请求同时在途,等待被互相掩盖.


两套配置模板

延迟优先(在线业务)

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),

// 不等待凑批: 低流量时单条即发, 高流量时自然攒批
kgo.ProducerLinger(0),
kgo.ProducerBatchMaxBytes(256 * 1024),

// lz4 CPU 开销低, 压缩不给延迟添乱
kgo.ProducerBatchCompression(kgo.Lz4Compression()),

// 可靠性不妥协(第 06 篇的组合)
kgo.RequiredAcks(kgo.AllISRAcks()),
)

吞吐优先(日志与数据管道)

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),

// 等 20ms 凑批: 用固定的延迟底噪换批量
kgo.ProducerLinger(20 * time.Millisecond),
kgo.ProducerBatchMaxBytes(1 << 20), // 1MB

// 压缩率优先, 带宽与磁盘双省
kgo.ProducerBatchCompression(kgo.ZstdCompressionLevel(3)),

// 缓冲按量放大, 吸收洪峰
kgo.MaxBufferedRecords(100_000),
)

两套模板的差异只有三个参数,但它们对应完全不同的业务假设.抄参数之前先回答:这套业务的 P99 预算是多少?是带宽瓶颈还是延迟瓶颈?


实测:先量后调

调参不能靠猜.官方压测工具可以快速跑对照:

docker compose exec kafka-0 bash

# 基线: Java 默认参数(16KB batch, linger 0, 无压缩)
kafka-producer-perf-test.sh \
--topic perf-test \
--num-records 1000000 \
--record-size 1024 \
--throughput -1 \
--producer-props bootstrap.servers=localhost:19092 acks=1

# 对照: 大批量 + linger + 压缩
kafka-producer-perf-test.sh \
--topic perf-test \
--num-records 1000000 \
--record-size 1024 \
--throughput -1 \
--producer-props bootstrap.servers=localhost:19092 \
acks=1 linger.ms=20 batch.size=262144 compression.type=lz4

--throughput -1 表示不限速.输出关注两个数:

1000000 records sent, 254432.50 records/sec (248.47 MB/sec), 12.30 ms avg latency, ...
▲ 吞吐 ▲ 延迟

调参循环:

1. 定基线: 当前参数下的吞吐与 P99
2. 单变量: 每轮只改一个参数, 一次改多个就分不清是谁起的作用
3. 复测: 同一压测条件重跑, 排除环境噪声
4. 上生产: 灰度发布, 看第 07 篇的 Producer 指标验证

生产验证时重点看两个指标:

  • batch-size-avg:远小于 batch.size 说明批根本没攒起来,linger 或流量不足,调大 batch.size 无意义
  • request-latency-avg 与 record-send-rate:延迟涨了吞吐没涨,说明调到了错误的方向

perf-test 的局限

kafka-producer-perf-test.sh 用的是 Java 客户端,默认值与 franz-go 不同(batch.size 16KB vs 1MB),压测结论要换算到自己的客户端上验证.
更好的做法:用它快速探索 Broker 与系统层的能力上限,用真实 Go 客户端跑最终配置的对照.


生产注意事项

先定位瓶颈,再调参

如果 Broker 的入站带宽、磁盘写入或网络线程先到上限(第 07 篇的 BytesInPerSec, RequestHandlerAvgIdlePercent),客户端怎么调都是徒劳.
顺序永远是:先量系统瓶颈在哪一段,再决定调客户端还是扩 Broker.

压缩 CPU 花在客户端

高吞吐 Producer 叠加 zstd 高压缩级别,Sender 线程可能先于网络成为瓶颈.
容器部署注意 CPU limit:调大压缩级别时同步评估 CPU 配额,别让业务与压缩抢核.

buffer 越大,故障丢失窗口越大

缓冲区里未发送的消息,进程崩溃即丢失(第 06 篇讲的 Producer 内存丢失场景).
buffer 调大的收益是吸收洪峰,代价是崩溃时丢得更多;关键业务配合 Outbox 模式,不要只靠调大 buffer.

调参纪律

一次一个变量,压测可复现,生产灰度验证,改完记录在案.
参数会随业务流量演进而漂移:今天的最优值,半年后可能因为流量翻倍而需要重测.


linger.ms=0 是不是每条消息一个网络请求?

不是.linger 0 的语义是"不等",同一个瞬间到达的多条消息仍然会进同一批,一起发走.
它保证的是不主动叠加等待,不阻止批量的自然形成.

batch.size 调到 1MB,低流量时会不会一直等满才发?

不会.batch.size 是上限而不是阈值,低流量下 linger 先到点,装了多少发多少.
真正的风险不是"等满",而是批太小导致逐条发送,这时候再看 batch-size-avg 指标决定是否调大 linger.

单条消息比 batch.size 还大怎么办?

不会被拆,单独成一批发送.消息大小的硬上限由 Broker 的 message.max.bytes(默认约 1MB)约束,与 batch.size 无关.
超过上限的消息会被 Broker 拒绝,需要调大 Broker 配置并评估全链路影响.

压缩会增加端到端延迟吗?

客户端多一次压缩 CPU 开销,但网络传输与 Broker 落盘的数据量都变小.
瓶颈在带宽或磁盘时,压缩让延迟反而更低;瓶颈在客户端 CPU 时,才是净负担.

acks=all 又想保吞吐,怎么两全?

批量摊薄 + 在途管道化(幂等模式下上限 5),把每次等待 ISR 的固定延迟藏到并发后面.
跨机房场景再叠加第 09 篇的机架感知部署,缩短 ISR 同步距离,减少每次等待的绝对值.

换客户端后参数怎么迁移?

语义一致但默认值不同:Java batch.size 16KB,franz-go 默认 1MB;Java buffer.memory 32MB,franz-go 按条数限流.
迁移时把两边的默认值都列出来,重新做一轮压测,不要假设同名参数行为相同.


快速回顾

  • 目标先行:延迟是 1) 等凑批 + 2) 网络 + 3) 等确认 + 4) 返回,只有 1) 可以主动伸缩换吞吐
  • 批量是第一杠杆:batch.size 定上限,linger.ms 定超时,谁先到谁触发;大 batch + linger 0 是通用起点
  • 压缩单位是 batch:批越大压缩越划算;lz4 默认首选,zstd 冲压缩率;压缩留在 Producer 端,Broker 保持 producer 模式
  • buffer.memory 是背压闸门:池满时 Produce 阻塞,把压力传回业务;容量按"峰值速率 × 可阻塞时长"估
  • 在途请求藏 RTT:幂等模式下 5 个在途,acks=all 也能高吞吐
  • 先量后调:定基线,单变量,复测,灰度;batch-size-avg 是验证批量是否真的形成的镜子

动手练习

  1. 建立基线:用 perf-test 默认参数压 100 万条 1KB 消息,记录吞吐与平均延迟,作为后续所有对照的基线.
  2. linger 扫描:固定其他参数,只改 linger.ms(0 / 5 / 20 / 50),记录四组吞吐与延迟,画出一条自己的取舍曲线.
  3. 压缩对比:对比 none / lz4 / zstd 三组,同时看吞吐、客户端 CPU(容器 top)与 Broker 端 BytesInPerSec(第 07 篇指标).
  4. 验证批大小:用第 07 篇的 batch-size-avg 指标(或 franz-go 的 kprom)确认 linger 与压缩组合下批到底攒了多大.
  5. 背压实验:Go 程序把 MaxBufferedRecords 调到很小,持续高速 Produce,测量 Produce() 的阻塞时长,体会背压传导.
  6. 两套模板对照:按"延迟优先"与"吞吐优先"两套模板各跑一遍压测,判断手头业务该选哪套,写出选择理由.

下一篇进入消费调优:fetch 参数、并发模型与背压处理.第 04 篇给过消费模型与并发模式的骨架,下篇讲 Go 消费者怎么落地 worker pool、慢消费怎么治理、fetch 参数怎么和业务处理速度匹配.