一句话总结
Consumer Group 是 Kafka 实现水平扩展消费的核心机制:同一 Group 内的 Consumer 瓜分 Partition,每个 Partition 同一时刻只被一个 Consumer 消费;不同 Group 之间互相独立,各自维护 offset.
理解 rebalance 和 offset 提交,就理解了消费端 80% 的坑.
消费模型:Pull 而非 Push
Kafka 消费者是拉模型(pull)——Consumer 主动向 Broker 发起 FetchRequest,拉取一批消息处理.
| 对比 | Push(RabbitMQ) | Pull(Kafka) |
|---|---|---|
| 谁主导 | Broker 推给 Consumer | Consumer 向 Broker 拉 |
| 速率控制 | Broker 控制,Consumer 被动接收 | Consumer 控制,按自身处理能力拉 |
| 背压 | 需要 prefetch 限制,否则压垮 Consumer | 天然背压:不拉就不来 |
| 空闲时 | 无消息时 Broker 不推 | Consumer 轮询空转(long poll 优化) |
类比 Go 的 channel:Pull 模型像 for msg := range ch 主动读;Push 模型像开个 goroutine 不断往你的 handler 塞数据,你来不及处理就溢出.
Long Polling:避免空转
Consumer 拉取时传入 fetch.min.bytes 和 fetch.max.wait.ms:
|
这样在消息稀少时不会疯狂空转,也不会延迟太高.
Consumer Group:协作消费的核心
为什么需要 Group
单个 Consumer 的消费速率有上限.如果一个 Topic 有 12 个 Partition,一个 Consumer 要顺序处理全部分区,吞吐就是瓶颈.
Consumer Group 允许多个 Consumer 实例分摊 Partition:
|
核心规则
| 规则 | 说明 |
|---|---|
| 一个 Partition 同一时刻只属于 Group 内的一个 Consumer | 保证分区内消息不被重复处理 |
| 一个 Consumer 可以消费多个 Partition | Consumer 数 < Partition 数时 |
| Consumer 数 > Partition 数时,多余的 Consumer 空闲 | 不会有两个 Consumer 分到同一 Partition |
| 不同 Group 之间完全独立 | 同一条消息可以被多个 Group 各消费一次 |
扩缩容与分区数的关系
|
关键结论
Consumer 实例数的有效上限 = Topic 的 Partition 数.
想提高消费并行度?增加 Partition 数(但要权衡之前讲的 Key 路由问题).
多 Group 独立消费:发布-订阅模式
不同 Group 对同一 Topic 的消费互不影响,各自维护 offset:
|
这就是 Kafka 实现"一条消息被多个系统消费"的方式——不需要消息复制,不需要 fanout exchange,只是不同 Group 各自读同一份日志的不同位置.
Rebalance:分区重分配
当 Group 成员发生变化时,Kafka 需要重新分配 Partition → Consumer 的映射,这个过程叫 Rebalance.
触发条件
| 事件 | 说明 |
|---|---|
| Consumer 加入 Group | 新实例启动,或原有实例重新连接 |
| Consumer 离开 Group | 正常关闭(Close())或心跳超时被踢出 |
| Topic 分区数变化 | 新增分区需要分配 |
| 订阅的 Topic 列表变化 | 正则订阅时新 Topic 匹配 |
Rebalance 的代价
Rebalance 期间所有 Consumer 停止消费(Eager 协议),这是最大的痛点:
|
典型 rebalance 耗时:秒级到十几秒(取决于 Group 大小和协调开销).对延迟敏感的服务来说,这是不可接受的"消费停顿".
两代 Rebalance 协议
| 协议 | 行为 | 问题 |
|---|---|---|
| Eager(老) | 所有 Consumer 先放弃全部 Partition,重新分配 | Stop-the-world,全员暂停 |
| Cooperative(新) | 只迁移需要变更的 Partition,其余继续消费 | 增量式,影响小 |
|
推荐:使用 Cooperative 协议
franz-go 默认使用 Cooperative rebalance(cooperative-sticky assignor).
如果你用 sarama,需要手动配置 Consumer.Group.Rebalance.GroupStrategies 为 CooperativeStickyBalancer.
分配策略(Assignor)
分配策略决定 Partition 怎么分给 Consumer:
| 策略 | 逻辑 | 特点 |
|---|---|---|
| Range | 按 Topic 内 Partition 连续分配 | 简单但容易不均匀 |
| RoundRobin | 全局轮转分配 | 均匀但 rebalance 变动大 |
| Sticky | 尽量保持原有分配不变 | 最小化 Partition 迁移 |
| CooperativeSticky | Sticky + Cooperative 协议 | 最优选,增量迁移 + 均匀 |
Offset 提交:记住消费到哪了
Consumer 需要告诉 Kafka"我处理到哪个 offset 了",这样 crash 后能从正确位置恢复.
存在哪
Offset 提交到 Kafka 内部的特殊 Topic:__consumer_offsets(50 个分区,compact 策略).
|
自动提交 vs 手动提交
| 模式 | 行为 | 风险 |
|---|---|---|
| 自动提交 | 每隔 auto.commit.interval.ms(默认 5s)自动提交当前位置 |
消息拉下来还没处理完就提交了 → crash 后丢消息 |
| 手动同步提交 | 处理完一批后调用 CommitSync() |
安全但阻塞,影响吞吐 |
| 手动异步提交 | 处理完调用 CommitAsync() |
不阻塞但失败时可能 offset 回退 |
自动提交的坑
自动提交的语义是"我拉到的最新 offset",不是"我处理完的最新 offset".
如果消息拉下来还在处理,自动提交已经把 offset 推进了,此时 crash → 那批消息就丢了.
生产环境强烈建议手动提交.
手动提交的正确姿势
|
提交粒度的权衡
| 策略 | 做法 | 效果 |
|---|---|---|
| 每批提交 | Poll 一批 → 处理 → 提交 | 平衡:最多重复处理一批 |
| 每条提交 | 处理一条 → 提交一条 | 最少重复,但提交开销大 |
| 定时提交 | 每 N 秒提交一次当前处理位置 | 吞吐好,但窗口内可能重复 |
生产环境通常用每批提交:一次 PollFetches 拿到的消息处理完后统一提交.如果 crash,最多重复这一批——配合业务幂等就够了.
Consumer 的心跳与会话
Consumer 通过心跳向 Group Coordinator(一个 Broker 角色)报活:
| 参数 | 含义 | 默认值 | 说明 |
|---|---|---|---|
heartbeat.interval.ms |
心跳发送间隔 | 3s | 应 < session.timeout 的 1/3 |
session.timeout.ms |
超时未心跳则踢出 | 45s | 太短 → 网络抖动误踢;太长 → 故障检测慢 |
max.poll.interval.ms |
两次 Poll 的最大间隔 | 5min | 处理太慢被判定"挂了" |
max.poll.interval.ms:最隐蔽的坑
心跳是后台线程发的,即使你的处理逻辑卡死了,心跳照样正常.所以 Kafka 用 max.poll.interval.ms 做第二道检测:
|
防御措施
- 控制每批拉取的消息量(
max.poll.records/ franz-go 的FetchMaxBytes) - 如果单条处理确实耗时,调大
max.poll.interval.ms - 或者:Pull 后快速放入本地 channel/worker pool,主循环立即 Poll 下一批
Go 服务中 Consumer 的粒度
一个常见疑问:Go 服务里"一个 Consumer"到底是什么?是一个进程?一个 goroutine?
答案:一个 kgo.Client 实例 = 一个 Consumer(Group 中的一个 member).它在 Kafka 协议层面持有唯一的 member ID,与 Group Coordinator 保持心跳.goroutine 只是进程内的并发单元,Kafka 完全感知不到.
|
扩展消费能力的正确方式
| 方式 | 做法 | Kafka 视角 | 效果 |
|---|---|---|---|
| 加 Pod 副本 | K8s scale replicas | Group 成员增加,触发 rebalance | 更多 Consumer 分摊 Partition |
| Pod 内加 goroutine | worker pool 处理消息 | 无感知,还是同一个 Consumer | 单 Consumer 处理能力提升 |
| 加 Partition 数 | kafka-topics --alter |
更多分区可分配 | 能容纳更多 Consumer 并行 |
为什么不在同一个进程里创建多个 Client
技术上可以在一个 Pod 里创建多个 kgo.Client 加入同一个 Group,但几乎没人这么做——因为瓶颈通常不在"拉取"而在"处理".一个 Client 拉到消息后开 N 个 goroutine 并行处理,效果一样,还省了多份心跳、TCP 连接和 rebalance 开销.
需要更多 Consumer?加 Pod 副本数,这是 K8s 时代的标准做法.
Go 中的并发消费模式
理解了"一个 Client = 一个 Consumer"后,接下来的问题是:单个 Consumer 内部怎么用 goroutine 提高处理吞吐?
模式一:每 Partition 一个 goroutine
|
|
模式二:Worker Pool(牺牲分区顺序换吞吐)
|
模式三:franz-go 的 BlockRebalanceOnPoll
franz-go 提供 BlockRebalanceOnPoll 选项:在你处理消息期间阻止 rebalance,直到你调用 AllowRebalance().这让你可以安全地用 goroutine 处理,不怕分区被抢走:
|
消费起始位置:从哪里开始读
新 Consumer Group 第一次消费一个 Topic 时,没有已提交的 offset,需要决定从哪开始:
auto.offset.reset |
行为 | 场景 |
|---|---|---|
latest |
从最新消息开始,忽略历史 | 只关心新事件(默认值) |
earliest |
从最早可用消息开始 | 需要消费全部历史数据 |
none |
找不到 offset 就报错 | 严格场景,不允许静默跳过 |
如果 Group 之前消费过(有提交记录),就从提交的 offset 继续.
auto.offset.reset 只影响两种情况:全新 Group、或已提交的 offset 对应的消息已过期被删除.
小结
| 概念 | 一句话 |
|---|---|
| Pull 模型 | Consumer 主动拉,天然背压,long poll 避免空转 |
| Consumer Group | 同 Group 内瓜分 Partition,不同 Group 独立消费 |
| Rebalance | 成员变化时重分配 Partition;Cooperative 协议减少停顿 |
| Offset 提交 | 手动提交 + 每批粒度,配合业务幂等 |
| 心跳 & Poll 间隔 | 两道防线:心跳检测网络存活,poll 间隔检测处理卡死 |
| 并发消费 | 按 Partition 分 goroutine 保序;Worker Pool 换吞吐 |
下一篇进入 Go 客户端实战:对比 sarama / confluent-kafka-go / franz-go 三大库的取舍,完整工程结构与优雅关停.