阶段二 · 架构原理与调优

消费调优:Fetch 参数,并发模型与背压

前置回顾

  • 第 04 篇给了消费模型与并发模式的骨架:Pull 拉取、每分区一个 goroutine、worker pool、BlockRebalanceOnPoll
  • 第 11 篇讲过慢处理如何撞上 max.poll.interval 并触发 rebalance,那是"坑",本篇是它的镜像面:调优
  • 第 12 篇调 Producer 端,消费端是镜像问题:一次拉回一批,逐条处理

一句话总结

消费端调优的目标不是"拉得快",而是"处理得完".
Fetch 参数决定一次拉多少,并发模型决定单位时间处理多少,背压决定拉回来的数据堆在哪里.三者错配的症状就三个:lag 涨、内存涨、rebalance 频繁.

一个慢消费的现场

订单处理服务:每条消息要调两次下游 RPC,各 50ms,单条处理耗时 100ms.Topic 6 分区,3 个消费者,按第 04 篇的"每分区一个 goroutine"消费.

处理能力算一下:

单分区串行处理: 1000ms / 100ms = 10 条/s
单消费者 2 分区: 20 条/s
整个组 3 消费者: 60 条/s

上游生产速率: 200 条/s

缺口 140 条/s,lag 单调上涨.不处理的结局是第 07 篇讲过的:lag 涨到 retention 之外,老消息被删,数据丢失.

这个现场里可调的只有三个维度:

维度 决定什么
Fetch 参数 一次拉回多少数据
并发模型 单个消费者的处理能力
背压 处理不过来时数据堆在哪

动手之前先把瓶颈分段,拉住"拉"和"处理"两个环节:

lag 持续增长, 先分清瓶颈在哪一段:

records-consumed-rate 上不去, CPU 不饱和
→ 拉得慢: 网络 / fetch 参数 / broker 端(第 07 篇)

consumed rate 正常, lag 不降, 内存上涨
→ 处理不完: 并发不够 / 单条太慢 / 下游是瓶颈

Fetch 参数:一次拉多少

第 04 篇讲过 Pull 模型与 long poll,这里把一轮 FetchRequest 的参数补全:

消费者 FetchRequest                          Broker
│ FetchMinBytes=1, FetchMaxWait=500ms │
├────────────────────────────────────────────►│
│ │ 数据不足 minBytes 时挂起等待
│ │ (最长等 MaxWait)
│◄────────────────────────────────────────────┤ 返回 ≤ FetchMaxBytes 的数据
│ (单个分区单次 ≤ FetchMaxPartitionBytes) │

参数对照(Java 名 / franz-go 选项 / 作用 / Java 默认值):

Java franz-go 作用 默认
fetch.min.bytes kgo.FetchMinBytes 攒够多少字节才返回 1
fetch.max.wait.ms kgo.FetchMaxWait 不够 minBytes 时最多等多久 500ms
fetch.max.bytes kgo.FetchMaxBytes 单次 fetch 总上限 50MB
max.partition.fetch.bytes kgo.FetchMaxPartitionBytes 单分区单次上限 1MB
max.poll.records kgo.MaxPollRecords 单次 Poll 返回条数上限 500

四个调优点:

  1. FetchMinBytes 调大(如 1MB):Broker 攒够再返回,少而大的响应,省往返.吞吐型消费的常用手段;延迟敏感场景保持小值
  2. FetchMaxPartitionBytes 是双刃剑:默认 1MB,大消息或大积压时单个分区要多次往返才能拉完;调大加快拉取,但单次 Poll 的内存占用也上涨
  3. FetchMaxBytes 是内存闸门:一次拿回的总量上限,处理慢的消费者别调太大
  4. MaxPollRecords 与处理速度有明确算式:
MaxPollRecords × 单条处理时间 < max.poll.interval.ms

例: 单条 100ms, max.poll.interval 默认 5min
→ 单批上限约 3000 条
→ MaxPollRecords 设 5000 且批批都要处理完, 必然超时被踢(第 04/11 篇)

这就是"fetch 参数与处理速度匹配"的全部含义:一次拉回的活,必须在 poll 超时前干完.

franz-go 的 MaxConcurrentFetches

默认每个 Broker 只有一个在途 fetch,逐 broker 串行.
调大 kgo.MaxConcurrentFetches 可让同一 Broker 上多个分区并行拉取:吞吐换内存,在途请求越多,单次 Poll 聚合的数据量越大.


Worker Pool:并发模型落地

第 04 篇给过三种模式的骨架,这里给出可落地的实现,重点是 offset 提交的正确处理.

方案一:每分区一个 goroutine,分区内保序

顺序敏感场景的标准解.每个分区一个 goroutine,提升单分区吞吐靠"批量化处理":

fetches.EachPartition(func(p kgo.FetchTopicPartition) {
wg.Add(1)
go func(topic string, partition int32, records []*kgo.Record) {
defer wg.Done()
// 关键优化: 慢的是 RPC, 不是计算 → 攒批调用
// 50ms 调 1 条 vs 50ms 调 100 条, 吞吐差两个数量级
processBatch(records)
}(p.Topic, p.Partition, p.Records)
})
wg.Wait()
client.CommitUncommittedOffsets(ctx)

并发度 = 分区数,分区内串行(保序),offset 提交简单:整批处理完统一提交.开头现场里的 100ms 单条 RPC,批量化后可能变成 50ms 处理 100 条.

方案二:共享 worker pool,吞吐优先

任何分区的消息交给任意 worker,并发度不受分区数限制,但分区内顺序丢失(第 02 篇:有序性只在分区内).实现的关键是 offset 提交不能"处理完一条提交一条":

Partition 0 各 offset 的处理状态:

offset 100 ✓ 101 ✓ 102 ✗(处理中) 103 ✓
▲
连续完成的最高水位是 101
→ 只能提交到 102, 不能提交到 104

若 103 被提前提交而 102 之后崩溃:
→ 重启后从 103 之后继续, 102 被跳过 → 丢消息

正确做法是维护"连续完成水位":

// 每个分区维护: 已完成集合 + 连续水位
type tracker struct {
mu sync.Mutex
done map[int64]bool // 已完成但尚未连成片
highmark int64 // 连续完成的最高 offset
}

func (t *tracker) mark(offset int64) int64 {
t.mu.Lock()
defer t.mu.Unlock()
t.done[offset] = true
for t.done[t.highmark+1] { // 从水位开始, 连续推进
delete(t.done, t.highmark+1)
t.highmark++
}
return t.highmark
}

提交时只提交各分区的 highmark + 1.复杂度不低,换来的是吞吐上限只受 worker 数约束.

方案三:有界队列 + 阻塞投递(最简背压)

records := make(chan *kgo.Record, 1000)   // 有界队列

// worker 侧
for i := 0; i < numWorkers; i++ {
go func() {
for r := range records {
process(r)
}
}()
}

// poll 侧
for {
fetches := client.PollFetches(ctx)
fetches.EachRecord(func(r *kgo.Record) {
records <- r // 队列满, 这里阻塞 → 拉取速度自动降到处理速度
})
// 注意: 此时不能立即提交 offset(队列里还有在途消息)
}

代价是 offset 提交又回到"有在途消息未完成"的问题,需要配合完成计数或水位追踪.前两个方案的处理能力不够时才上这一层.


背压处理:数据堆在哪

背压 = 水库的闸门:上游(生产)来水不停,下游(处理)放水有限,水库(消费者缓冲)蓄满时必须关闸,否则水漫堤(内存爆).
关闸的时刻,水留在上游(Broker 的 lag 里),数据是安全的;不关闸,数据涌进本地内存,进程先死.

三种做法的对比:

做法 拉取速度 数据堆在哪 风险
不设限 全速拉 消费者内存 OOM,进程被杀,再触发 rebalance
阻塞投递(有界队列) 自动降到处理速度 Broker(lag 涨) 无 OOM 风险;要处理 offset 提交
暂停/恢复拉取 按水位手动控制 Broker 需要水位监控代码

franz-go 提供了显式的暂停接口,适合"本地队列 + 水位"的精细控制:

// 本地积压超过高水位: 暂停拉取
if localQueue.Len() > highWatermark {
client.PauseFetchTopics("order-events")
}

// 消化到低水位: 恢复
if localQueue.Len() < lowWatermark {
client.ResumeFetchTopics("order-events")
}

需要更细粒度时还有分区级的 PauseFetchPartitions / ResumeFetchPartitions.

背压的判断标准

问一个问题:处理不过来时,数据堆在 broker(安全,只是 lag 涨)还是堆在消费者内存(危险,会 OOM)?
有背压 = 堆在 broker.没有背压 = 内存无界增长,进程被杀,触发 rebalance,恶性循环.


慢消费治理:按顺序排查

1. 量化: lag 增速 = 生产速率 - 消费速率; 清空时间 = lag / 净消费速率
2. 分段: 瓶颈在"拉"还是"处理"(本文开头的判断树)
3. 处理侧按性价比选手段(见下表)
4. 检查连锁反应: max.poll.interval 是否被突破(触发 rebalance)

处理侧的手段,按投入产出排序:

手段 典型收益 代价
批量化(攒批 RPC、批量写库) 数量级 需要改处理逻辑
优化单条(缓存、减少 RPC) 数倍 工程成本
Worker Pool 并发 × worker 数 分区内顺序丢失
加消费者实例 × 实例数,上限 = 分区数 触发 rebalance(第 11 篇)
加分区数 突破并行度上限 第 03 篇的 Key 路由问题

两个容易忽略的点:

  • 加消费者有上限:消费者数超过分区数,多出来的实例拿不到分区(第 04 篇),再加只是浪费资源
  • 慢消费与 rebalance 是连锁的:单批处理时间超过 max.poll.interval.ms,分区被夺走,进入第 11 篇的排查清单.止血手段之一就是把 MaxPollRecords 调小,让单批更快回到 Poll

动手验证

# 1. 制造消费压力: perf-test 持续生产(第 12 篇的命令)
kafka-producer-perf-test.sh --topic order-events \
--num-records 500000 --record-size 512 --throughput 2000 \
--producer-props bootstrap.servers=localhost:19092

# 2. 观察 lag(第 07 篇的命令)
kafka-consumer-groups.sh --describe --group order-processor \
--bootstrap-server localhost:19092

代码侧对照实验:

  1. 处理函数 sleep 100ms, 分别用 1 / 4 / 16 个 worker, 观察 lag 曲线
  2. MaxPollRecords 设成让单批必然超时, 观察分区被夺走与 rebalance(接第 11 篇)
  3. 有界队列调小(如 100), 观察消费速率自动降到处理速度、内存保持稳定
  4. 把逐条 RPC 改成攒批 RPC, 对比同样并发下 lag 清空时间的数量级变化

生产注意事项

不要在 Poll 循环里做阻塞操作

第 11 篇从 rebalance 角度讲过这个坑,从吞吐角度再看一遍:Poll 循环里的任何阻塞(同步 RPC、大文件写、锁等待)会同时压住"拉取"和"提交".
固定模式: Poll → 投递队列/goroutine → 立即回到 Poll.

FetchMaxBytes 是内存闸门

一次 Poll 拿回的数据真实占用消费者内存.FetchMaxBytes 50MB 加多个在途 fetch(每 broker 一个),峰值内存可以轻松到几百 MB.
容器内存紧张时先算这笔账,再决定 fetch 参数的尺寸.

顺序与并发的取舍要显式做

Worker Pool 换吞吐,代价是分区内消息乱序(第 02 篇:有序性只在分区内).
需要顺序的业务退回"每分区串行 + 批量化";不需要顺序才上共享池,不要默认选择"更快"的那个.

消费端稳态配方

每分区 goroutine + 分区内批量化 + 有界缓冲 + 水位暂停 + 手动提交:
顺序不丢、吞吐可扩、内存有界、重复窗口可控.


加了消费者 lag 还是不降,为什么?

按顺序查三件事:消费者数是否已等于分区数(第 04 篇的上限);瓶颈是否在共享资源(数据库连接池、下游限流);rebalance 是否频繁发生(第 11 篇的排查清单).
三者都不是,瓶颈可能在"拉":检查 broker 端指标(第 07 篇)与网络.

FetchMinBytes 调大,会不会增加延迟?

会.流量小时,broker 要攒够 minBytes 才返回,不够就等到 MaxWait(如 500ms).
延迟敏感的消费者保持 minBytes=1;吞吐型消费者调大,用"延迟增加几十毫秒"换"请求数下降一个数量级".

worker pool 下 offset 怎么提交才对?

不能"处理完一条提交一条":乱序完成时,提交大的会跳过未完成的小 offset(丢消息).
要么退回每分区串行(天然有序);要么维护"连续完成水位",只提交水位(本节方案二的代码).提交语义仍是 at-least-once:水位之前保证处理过,水位附近可能重复.

消费端能不能做到不丢不重?

单靠客户端配置做不到,这是第 06 篇的结论.
消费端 EOS 要靠事务(consume-transform-produce 场景),下一篇拆开讲.

消费者内存占用怎么估?

粗略上界:在途 fetch 数 × FetchMaxBytes + 本地缓冲队列 + 处理中的消息.
默认的"Poll 循环不阻塞 + 无界处理"下这个值没有上界,所以背压不是优化项,是正确性问题.


快速回顾

  • 目标:消费调优是"处理得完",不是"拉得快";三种症状,lag 涨、内存涨、rebalance 频繁
  • Fetch 参数:MinBytes/MaxWait 控制长轮询,MaxBytes/MaxPartitionBytes 控制单次体量,MaxPollRecords 与处理速度有明确算式
  • Worker Pool:分区内串行保序;共享池换吞吐,但要用"连续完成水位"才能正确提交 offset
  • 背压:数据堆在 broker(lag)而不是消费者内存(OOM);阻塞投递或 Pause/Resume 都是合法实现
  • 慢消费治理:先分段定瓶颈,处理侧按"批量化 > 优化 > 并发 > 加实例 > 加分区"投入
  • 连锁反应:慢处理会撞上 max.poll.interval 触发 rebalance,治理时连着看第 11 篇的清单

动手练习

  1. 复现慢消费:写一个 sleep 100ms/条 的消费者,用 perf-test 打 2000 条/s,观察 lag 单调上涨.
  2. 并发对照:同一处理逻辑,1 / 4 / 16 个 worker 各跑一遍,画出 lag 与内存两条曲线.
  3. 撞超时:把 MaxPollRecords 调到让单批必然超过 max.poll.interval,观察分区被夺走与 rebalance 日志(接第 11 篇).
  4. 批量化改造:把逐条 RPC 改成攒 100 条一次 RPC,对比改造前后的 lag 清空时间.
  5. 背压实验:有界队列从 10000 调到 100,观察消费速率自动降到处理速度,而内存保持稳定.
  6. fetch 参数实验:FetchMinBytes 从 1 调到 1MB,对比请求数与端到端延迟的变化.

下一篇是本阶段最后一篇:事务与 Exactly-Once.消费端"不丢不重"的完整答案在那里,跨分区原子写入、事务型生产者与消费者,以及 Kafka 事务的能力边界.