阶段一 · 基础与核心模型

消费者与 Consumer Group

一句话总结

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:

"Broker,给我拉数据,但如果不够 1KB,你最多等 500ms 再回复我"

这样在消息稀少时不会疯狂空转,也不会延迟太高.


Consumer Group:协作消费的核心

为什么需要 Group

单个 Consumer 的消费速率有上限.如果一个 Topic 有 12 个 Partition,一个 Consumer 要顺序处理全部分区,吞吐就是瓶颈.

Consumer Group 允许多个 Consumer 实例分摊 Partition:

Topic: order-events (6 partitions)

Consumer Group: "order-processor"
├── Consumer A → P0, P1
├── Consumer B → P2, P3
└── Consumer C → P4, P5

每个 Consumer 处理 2 个 Partition,吞吐翻 3 倍

核心规则

规则 说明
一个 Partition 同一时刻只属于 Group 内的一个 Consumer 保证分区内消息不被重复处理
一个 Consumer 可以消费多个 Partition Consumer 数 < Partition 数时
Consumer 数 > Partition 数时,多余的 Consumer 空闲 不会有两个 Consumer 分到同一 Partition
不同 Group 之间完全独立 同一条消息可以被多个 Group 各消费一次

扩缩容与分区数的关系

6 Partitions, Group 内 Consumer 数量变化:

1 个 Consumer: C1 → [P0, P1, P2, P3, P4, P5] (独扛全部)
3 个 Consumer: C1 → [P0, P1], C2 → [P2, P3], C3 → [P4, P5]
6 个 Consumer: C1→P0, C2→P1, C3→P2, C4→P3, C5→P4, C6→P5
7 个 Consumer: C7 空闲,分不到 Partition ← 浪费!

关键结论

Consumer 实例数的有效上限 = Topic 的 Partition 数.
想提高消费并行度?增加 Partition 数(但要权衡之前讲的 Key 路由问题).


多 Group 独立消费:发布-订阅模式

不同 Group 对同一 Topic 的消费互不影响,各自维护 offset:

Topic: order-events

Group "order-processor" → 处理订单业务逻辑
offset: P0=1520, P1=980, P2=2100

Group "analytics-pipeline" → 写入数据仓库
offset: P0=1200, P1=980, P2=1800 (消费更慢,offset 落后)

Group "audit-logger" → 写入审计日志
offset: P0=1520, P1=980, P2=2100 (和第一个一样快)

这就是 Kafka 实现"一条消息被多个系统消费"的方式——不需要消息复制,不需要 fanout exchange,只是不同 Group 各自读同一份日志的不同位置.


Rebalance:分区重分配

当 Group 成员发生变化时,Kafka 需要重新分配 Partition → Consumer 的映射,这个过程叫 Rebalance.

触发条件

事件 说明
Consumer 加入 Group 新实例启动,或原有实例重新连接
Consumer 离开 Group 正常关闭(Close())或心跳超时被踢出
Topic 分区数变化 新增分区需要分配
订阅的 Topic 列表变化 正则订阅时新 Topic 匹配

Rebalance 的代价

Rebalance 期间所有 Consumer 停止消费(Eager 协议),这是最大的痛点:

时间线:
──正常消费──┤ rebalance 开始 ├──暂停──┤ rebalance 完成 ├──恢复消费──
│ │
└── 这段时间无消息被处理 ──┘

典型 rebalance 耗时:秒级到十几秒(取决于 Group 大小和协调开销).对延迟敏感的服务来说,这是不可接受的"消费停顿".

两代 Rebalance 协议

协议 行为 问题
Eager(老) 所有 Consumer 先放弃全部 Partition,重新分配 Stop-the-world,全员暂停
Cooperative(新) 只迁移需要变更的 Partition,其余继续消费 增量式,影响小
Eager Rebalance (3 Consumer, 6 Partition):
C2 挂了 → C1,C3 同时放弃全部 Partition → 重新分配 → 恢复
影响:6 个 Partition 全部暂停

Cooperative Rebalance:
C2 挂了 → 只有 C2 原来负责的 P2,P3 需要迁移
→ C1,C3 继续处理自己的 P0,P1,P4,P5
→ P2,P3 被重新分配给 C1 或 C3
影响:只有 2 个 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 策略).

key: (group_id, topic, partition)
value: offset + metadata + timestamp

自动提交 vs 手动提交

模式 行为 风险
自动提交 每隔 auto.commit.interval.ms(默认 5s)自动提交当前位置 消息拉下来还没处理完就提交了 → crash 后丢消息
手动同步提交 处理完一批后调用 CommitSync() 安全但阻塞,影响吞吐
手动异步提交 处理完调用 CommitAsync() 不阻塞但失败时可能 offset 回退

自动提交的坑

自动提交的语义是"我拉到的最新 offset",不是"我处理完的最新 offset".
如果消息拉下来还在处理,自动提交已经把 offset 推进了,此时 crash → 那批消息就丢了.
生产环境强烈建议手动提交.

手动提交的正确姿势

// franz-go 手动提交示例
client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.ConsumerGroup("order-processor"),
kgo.ConsumeTopics("order-events"),
kgo.DisableAutoCommit(), // 关闭自动提交
)
defer client.Close()

for {
fetches := client.PollFetches(context.Background())
if errs := fetches.Errors(); len(errs) > 0 {
for _, e := range errs {
log.Error("fetch error", "err", e.Err, "topic", e.Topic, "partition", e.Partition)
}
}

fetches.EachRecord(func(r *kgo.Record) {
// 处理消息
processOrder(r.Value)
})

// 处理完这批再提交
if err := client.CommitUncommittedOffsets(context.Background()); err != nil {
log.Error("commit failed", "err", err)
}
}

提交粒度的权衡

策略 做法 效果
每批提交 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 做第二道检测:

PollFetches() → 处理消息(如果处理了 6 分钟)→ PollFetches()
↑
超过 max.poll.interval(默认 5min)
→ Coordinator 认为你挂了
→ 触发 Rebalance,你的 Partition 被抢走
→ 你还在处理上一批消息...
→ 处理完提交 offset → 失败(已不属于你)

防御措施

  1. 控制每批拉取的消息量(max.poll.records / franz-go 的 FetchMaxBytes)
  2. 如果单条处理确实耗时,调大 max.poll.interval.ms
  3. 或者:Pull 后快速放入本地 channel/worker pool,主循环立即 Poll 下一批

Go 服务中 Consumer 的粒度

一个常见疑问:Go 服务里"一个 Consumer"到底是什么?是一个进程?一个 goroutine?

答案:一个 kgo.Client 实例 = 一个 Consumer(Group 中的一个 member).它在 Kafka 协议层面持有唯一的 member ID,与 Group Coordinator 保持心跳.goroutine 只是进程内的并发单元,Kafka 完全感知不到.

典型部署:K8s 3 个 Pod,Topic 6 Partitions

Pod 1 (kgo.Client) → Consumer member-1 → P0, P1
Pod 2 (kgo.Client) → Consumer member-2 → P2, P3
Pod 3 (kgo.Client) → Consumer member-3 → P4, P5

每个 Pod 内部:
┌─────────────────────────────────────────────┐
│ kgo.Client (1 个 Consumer,持有 P0, P1) │
│ │ │
│ ├── goroutine: poll loop │
│ ├── goroutine: heartbeat │
│ ├── goroutine: worker-1 (处理 P0) │ ← Kafka 不知道
│ └── goroutine: worker-2 (处理 P1) │ ← 只是进程内并发
└─────────────────────────────────────────────┘

扩展消费能力的正确方式

方式 做法 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

PollFetches()
├── Partition 0 的消息 → goroutine 0 处理
├── Partition 1 的消息 → goroutine 1 处理
└── Partition 2 的消息 → goroutine 2 处理
每个 goroutine 内部串行 → 保证分区内有序
fetches.EachPartition(func(p kgo.FetchTopicPartition) {
go func(records []*kgo.Record) {
for _, r := range records {
process(r)
}
// 这批 partition 处理完,标记可提交
}(p.Records)
})

模式二:Worker Pool(牺牲分区顺序换吞吐)

PollFetches() → 所有消息塞入 channel → N 个 worker 并行处理
优点:吞吐最大化
缺点:同一 Partition 的消息可能乱序
适用:业务不要求严格顺序,或天然幂等

模式三:franz-go 的 BlockRebalanceOnPoll

franz-go 提供 BlockRebalanceOnPoll 选项:在你处理消息期间阻止 rebalance,直到你调用 AllowRebalance().这让你可以安全地用 goroutine 处理,不怕分区被抢走:

client, _ := kgo.NewClient(
// ...
kgo.BlockRebalanceOnPoll(),
)

for {
fetches := client.PollFetches(ctx)
// 在这里安全地并发处理,rebalance 被阻塞
processInParallel(fetches)
client.AllowRebalance()
client.CommitUncommittedOffsets(ctx)
}

消费起始位置:从哪里开始读

新 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 三大库的取舍,完整工程结构与优雅关停.