Go 基础体系 · 第 69/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go Kafka 完整指南:分区、消费者组、Offset 与 franz-go
本文以 Go 1.26.4、Apache Kafka 4.0.x(KRaft 模式)和 github.com/twmb/franz-go v1.19.x 为基线。Kafka 4.0 不再以 ZooKeeper 作为生产元数据模式;升级存量集群必须遵循实际来源版本的迁移路径。客户端与 Broker 的协议能力按 ApiVersions 协商,但锁定版本并做故障测试仍是发布前提。
Kafka 是分区追加日志:生产者决定记录进入哪个 partition,leader 按 offset 追加,副本复制,消费者以 group 和 offset 跟踪进度。它擅长高吞吐、长期保留和重放,却不会替应用解决业务数据库双写、跨分区全局顺序、毒消息和任意外部副作用的 exactly-once。
1. KRaft 架构与元数据职责
Broker 保存 topic partition 的日志并服务 Produce/Fetch。Controller quorum 使用 Raft 管理 topic、partition、leader、ISR、配置和 ACL 等集群元数据。生产可采用独立 controller 节点,避免数据流量与 controller 选举争抢资源;小环境可 combined,但故障域和容量边界不同。
客户端用 bootstrap brokers 取得 metadata,之后直接连接目标 partition leader。bootstrap 列表不是所有请求的代理,也不要求列出全部 Broker。Controller 不承载每条业务记录。一次 leader 变化后客户端刷新 metadata 并重试,因此监控连接成功不等于目标 partition 可写。
Producer -> metadata -> partition leader -> follower replicas
Consumer group -> group coordinator -> assignment -> partition leaders
Controllers -> __cluster_metadata quorum
2. Topic、Partition、Segment 与 Record Batch
Topic 是逻辑日志,partition 是有序追加和并行度的单位。每条 record 有 key、value、headers、timestamp 和写入后 offset;多个 record 以 batch 编码、压缩和校验。磁盘日志按 segment 切分,索引把 offset/time 映射到文件位置。保留策略按时间或大小删除旧 segment,compact topic 则按 key 最终保留较新值和 tombstone 语义。
offset 只在 partition 内有意义,不是全局时间。消费者不能用 offset 比较不同 partition 的先后。partition 数决定消费者组最大有效并行度之一,后期增加 partition 会改变默认 key 哈希映射,可能破坏同一 key 的连续顺序并增加文件、复制和选举成本。
3. Replica、Leader、ISR 与提交边界
每个 partition 有一个 leader 和若干 follower。生产、普通消费面向 leader;follower 拉取 leader 日志。ISR 是当前跟得上并有资格参与可靠确认的副本集合。复制因子 3 表示三份副本,但是否安全还取决于副本分布、ISR 和确认配置。
生产端 acks=all 要求当前 ISR 满足确认;Broker 的 min.insync.replicas=2 可在三副本下要求至少两个同步副本,否则拒绝写入。允许 unclean leader election 可能用落后副本恢复可用性却丢已确认数据,关键 topic 通常禁用。机架感知应把副本分布到不同故障域,不能三份都落同一宿主机或机房。
default.replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false
controller.quorum.voters=1@ctrl-1:9093,2@ctrl-2:9093,3@ctrl-3:9093
4. 一次 Produce 的生命周期
应用创建 record;客户端按 key 或分区器选择 partition,把记录放入有界 batch;达到字节阈值或 linger 后发送给 leader。Broker 校验权限、epoch 和序号,追加本地日志,按 acks 等待复制条件,再返回 base offset 或错误。客户端可能自动刷新 metadata、退避并重试。
压缩发生在 batch 级,批次越充分压缩率通常越好,但 linger 增加首条等待。acks=0 不返回 Broker 结果,接近至多一次;acks=1 只等待 leader,本地追加后 leader 未复制即故障可能丢;acks=all 配合 ISR 配置提供更强持久性,但响应丢失仍让客户端不知道结果。
请求超时不等于 Broker 未写。生产者重试必须使用幂等生产能力和稳定业务 event ID;即使 Kafka 日志避免同一 producer session 的协议级重复,业务在新 session、跨集群复制或下游处理仍可能看到重复。
5. franz-go 客户端创建与关闭
kgo.Client 并发安全,应长期复用。默认能力会随库版本变化,生产配置需显式表达可靠性、批量、超时和认证意图。SeedBrokers 只用于引导;地址应使用 Broker 对客户端可达的 advertised listeners。
func newKafkaClient(brokers []string, group string) (*kgo.Client, error) {
client, err := kgo.NewClient(
kgo.SeedBrokers(brokers...),
kgo.ClientID("article-indexer"),
kgo.ConsumerGroup(group),
kgo.ConsumeTopics("article.events"),
kgo.RequiredAcks(kgo.AllISRAcks()),
kgo.DisableAutoCommit(),
kgo.FetchMaxBytes(32<<20),
kgo.FetchMaxPartitionBytes(4<<20),
)
if err != nil {
return nil, fmt.Errorf("create kafka client: %w", err)
}
return client, nil
}
启动依赖应调用带 deadline 的 Ping 并检查目标 topic 授权,而不是只看构造成功。关闭生产者前用有界 context Flush,再 Close;Flush 期间不能让其他 goroutine 无界继续 Produce。消费者关停要停止 poll、等待处理完成并提交符合语义的 offset,然后离组。
6. 幂等生产者、重试和错误分类
幂等 producer 为每个 partition 维护 producer ID、epoch 和递增 sequence,Broker 用它过滤重试产生的重复并保持允许的在途顺序。franz-go 默认启用幂等写入时,不要随意用配置破坏其约束。幂等范围是 producer session 到 Kafka partition,不是业务数据库与 Kafka 的原子事务,也不永久按 MessageId 去重。
func publishArticle(ctx context.Context, client *kgo.Client, event ArticleEvent) error {
payload, err := json.Marshal(event)
if err != nil {
return fmt.Errorf("encode article event: %w", err)
}
record := &kgo.Record{
Topic: "article.events",
Key: []byte(event.ArticleID),
Value: payload,
Headers: []kgo.RecordHeader{
{Key: "event-id", Value: []byte(event.EventID)},
{Key: "schema-version", Value: []byte("1")},
},
}
result := client.ProduceSync(ctx, record)
if err := result.FirstErr(); err != nil {
return fmt.Errorf("produce article event %q: %w", event.EventID, err)
}
return nil
}
异步 Produce callback 不能忽略错误,也不能引用随后被调用方修改的 byte slice。达到本地缓冲限制时客户端会形成背压或返回错误;应用还要设置总 deadline。永久错误包括无权限、非法 topic 或超大消息;leader 变化、配额和暂时网络错误可能重试;超时后结果可能未知。无限重试会阻塞关停并扩大事故。
7. 分区选择、顺序和热点 key
带 key 的记录由分区器稳定映射,使同一 key 在 partition 数不变时进入同一 partition;无 key 记录通常按粘性策略形成批次。以 article ID/order ID 为 key 可得到实体内顺序。热点 key 会让一个 partition 成为瓶颈,随机加盐虽均衡吞吐,却失去直接的 key 内顺序,需要下游按版本归并。
Kafka 只保证同一 partition 的日志顺序。生产失败重试、多个独立 producer 同写一个 key 和消费者并行 handler 都可能改变业务完成顺序。事件携带聚合版本,消费者用“只接受期望版本/不回退版本”保护状态。增加 partition 前应评估 key 映射变化,必要时切新 topic 或在迁移期路由新老版本。
8. Consumer Group、Fetch 与再均衡生命周期
Group coordinator 管理成员和已提交 offset,group leader/协议计算 assignment。一个 partition 在一个 group 的稳定 generation 内只分配给一个成员;不同 group 可独立读取全部记录。消费者 Poll 获取 fetch 结果,也驱动心跳、再均衡回调和内部状态,不能长时间停止调用而不理解客户端的 poll 模型。
成员加入、退出、订阅变化或 partition 增加会触发再均衡。eager 策略先撤销全部分区;cooperative sticky 可渐进迁移,减少停顿,但应用仍须在 revoked 时停止旧分区工作并提交安全进度。静态成员能降低短重启抖动,却可能让失效成员在 session timeout 前占位。
9. Offset 提交与至少一次处理
提交 offset N 表示下次从 N 开始,通常应在成功处理 record offset N-1 后提交。先提交再处理会在崩溃时丢业务;先处理再提交会在业务成功、提交前崩溃时重复,因此常见语义是至少一次。
func consume(ctx context.Context, client *kgo.Client, handle func(context.Context, *kgo.Record) error) error {
for {
fetches := client.PollFetches(ctx)
if err := fetches.Err(); err != nil {
return fmt.Errorf("poll kafka: %w", err)
}
var processErr error
fetches.EachRecord(func(record *kgo.Record) {
if processErr != nil {
return
}
processErr = handle(ctx, record)
})
if processErr != nil {
return fmt.Errorf("process kafka record: %w", processErr)
}
if err := client.CommitRecords(ctx, fetches.Records()...); err != nil {
return fmt.Errorf("commit kafka records: %w", err)
}
}
}
该示例串行处理一个 poll 批次,易说明但吞吐有限。并行处理必须按 partition 跟踪“连续完成水位”,不能看到 offset 12 完成就提交 13,而 offset 10 仍失败。提交调用的同步/异步与再均衡行为必须按 franz-go 锁定版本验证。
10. 幂等消费、Outbox 与双写边界
外部数据库消费者把 event ID 登记和业务更新放在同一数据库事务,唯一冲突作为已处理。业务 API 生产事件时,将状态和 Outbox 行放在同一事务;relay 读 Outbox 发布 Kafka,确认后标记。relay 崩溃会重复发布,但下游幂等吸收,避免“数据库成功而事件永久缺失”。
不要用“先写数据库再 Produce,失败就重试”冒充原子性,因为进程可在两步间崩溃。CDC 可读取数据库日志替代轮询 Outbox,但 schema 演进、事务边界、重放和删除仍需治理。业务唯一键比短期缓存可靠;缓存可减负但不能成为唯一正确性来源。
11. Kafka 事务与 exactly-once 的真实范围
事务 producer 用稳定 transactional ID 初始化 epoch,把多个 partition 的记录和消费 offset 原子提交到 Kafka。下游配置 read_committed 才会跳过 aborted 数据。典型 consume-transform-produce 流程可在 Kafka 内实现“输入 offset 与输出记录同事务”,实例崩溃后旧 producer 被 fencing。
err := client.BeginTransaction()
if err != nil {
return fmt.Errorf("begin kafka transaction: %w", err)
}
// Produce 转换后的记录,并检查 callback/flush 结果。
offsets := client.MarkedOffsets()
if err := client.SendOffsetsToTransaction(ctx, offsets, groupMetadata); err != nil {
_ = client.AbortBufferedRecords(ctx)
return fmt.Errorf("send transaction offsets: %w", err)
}
if err := client.EndTransaction(ctx, kgo.TryCommit); err != nil {
return fmt.Errorf("commit kafka transaction: %w", err)
}
具体 API 需按 v1.19.x 编译验证,生产代码必须完整处理 Produce 结果和 abort。Kafka 事务不能把 MySQL 更新、HTTP 支付或邮件发送纳入同一原子提交;这类副作用仍需幂等/Outbox。事务增加协调、延迟和运维成本,只有确实需要 Kafka-to-Kafka 原子流时使用。
12. 重试 Topic、死信与背压
在原 partition 遇到永久坏记录时停住可保持顺序,却阻塞后续;跳过并发到 retry topic 可恢复吞吐,却破坏该 key 的处理顺序。系统必须明确选择。常见拓扑按延迟创建 retry topic,最终写 DLT;记录原 topic/partition/offset、event ID、attempt 和错误分类,避免把凭证与敏感 body 放 header。
消费者背压优先暂停 partition 或降低 poll 后的在途量,而不是继续 fetch 到无界 channel。FetchMaxBytes、partition fetch 上限、客户端最大缓存、worker 数和下游连接池共同构成预算。长任务可 Pause 后继续 Poll 维持组成员状态,或把工作拆为可快速持久化的内部任务;必须保证恢复、撤销和提交顺序。
生产侧本地缓冲满代表容量不足,应阻塞到 context、拒绝或落 Outbox。Broker quota/throttle 是保护信号,不应由更多并发重试抵消。监控 lag 的增长速度和最老事件年龄,单个 offset lag 在消息大小差异大时不能代表剩余工作量。
13. 故障恢复、保留和灾备
Broker 故障触发 leader 选举,客户端收到可重试错误并刷新 metadata。若 ISR 不满足 min.insync.replicas,可靠配置会拒绝写入而非降级丢数据。恢复 Broker 追赶期间磁盘和网络压力升高,应限制副本恢复流量并观察 under-replicated partitions。
retention 是在线日志生命周期,不是备份。误删 topic、错误 ACL 和应用批量写坏数据可能快速扩散。跨集群复制可用于灾备和地域分发,但 offset、transactional ID、ACL 和切换语义需演练;不能仅看到目标有 topic 就认为可无损切换。恢复演练验证 RPO/RTO、consumer 起点、重复处理和回切。
compact topic 只保证对 key 的最终保留趋势,清理由后台完成;墓碑也受保留期影响。依赖它重建状态时要测试从最早可用 offset 全量回放,并确保 key/schema 兼容。
14. 诊断、监控和运维命令
Broker 关注 request P99、网络/请求队列、磁盘使用与 I/O、under-replicated/offline partitions、ISR shrink、controller 活跃、KRaft quorum、quota throttle 和日志清理。Producer 关注 record error/retry、batch size、compression ratio、buffer wait 和 request latency;Consumer 关注按 partition lag、消费速率、rebalance、commit error 和处理耗时。
kafka-topics.sh --bootstrap-server kafka-1:9092 --describe --topic article.events
kafka-consumer-groups.sh --bootstrap-server kafka-1:9092 \
--describe --group article-indexer
kafka-metadata-quorum.sh --bootstrap-server kafka-1:9092 describe --status
kafka-configs.sh --bootstrap-server kafka-1:9092 \
--entity-type topics --entity-name article.events --describe
诊断单个事件时记录 topic、partition、offset、event ID、group 和 generation,不把它们都做指标 label。lag 为零但业务仍失败时检查 handler 成功条件和 DLT;lag 高但吞吐正常可能只是上游突发。时间戳受生产者时钟和 topic 配置影响,不用它替代严格因果顺序。
15. 测试、性能与故障注入
无集群单元测试覆盖 key 分区策略、事件编码、重试分类、连续 offset 水位、幂等仓储和 DLT 决策。协议集成测试固定 Kafka 4.0 镜像,验证 acks、压缩、超大记录、consumer group、再均衡、手动提交、read_committed 和 ACL。多 Broker 环境注入 leader 停止、ISR 缩小、网络延迟、磁盘水位和 controller 切换。
gofmt -w ./verify/kafka/*.go
go test ./verify/kafka/...
go test -race ./verify/kafka/...
kafka-producer-perf-test.sh --topic article.events --num-records 1000000 \
--record-size 1024 --throughput -1 \
--producer-props bootstrap.servers=kafka-1:9092 acks=all compression.type=zstd
基准使用真实 key 倾斜、record/header 大小、压缩、partition 和副本数,记录 MB/s、records/s、P50/P99、CPU、磁盘、网络、batch fill 与恢复追赶时间。峰值吞吐之外测试持续保留、滚动升级、积压后追赶及下游变慢;只测 producer 会遗漏消费者和磁盘清理的真实瓶颈。
16. 安全、部署与选型边界
生产 listener 使用 TLS,跨信任域使用 SASL/SCRAM、OAuth 或 mTLS,并通过 ACL 限制 topic、group、transactional ID 和集群操作。禁用明文公网端口,保护 controller listener。凭证放 secret 并轮换;日志不输出 SASL 配置或完整 payload。限制最大消息、请求和连接配额,schema registry/契约校验防止不兼容事件污染长期日志。
Broker 使用独立持久盘和稳定网络,按机架/可用区分布副本,给重建、retention 和 compaction 留空间。分区不是越多越好:它增加 fd、内存、复制、选举和客户端元数据成本。滚动升级前检查支持的协议与 metadata version,升级后再按官方步骤提升功能级别,保留回滚窗口。
Kafka 适合高吞吐事件流、CDC、审计日志、长期保留、多消费者独立重放和流处理。需要复杂按内容路由、逐消息 Ack、短任务优先级时 RabbitMQ 更直接;需要面向交易的延迟/事务回查生态可评估 RocketMQ;低量同步调用不应绕成消息流。最终选择应以保留与重放模型、顺序粒度、故障语义、运维成本和团队验证能力为依据,而不是仅看基准数字。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go RabbitMQ 实战:Exchange、Queue、确认、重试与死信
- 下一篇:Go RocketMQ 实战:Producer、Consumer、顺序与延迟消息
- 延伸:Go 可靠消息统一设计:Outbox、幂等、重试、顺序与死信
- 延伸:Go ClickHouse 实战:列式建模、批量写入与分析查询
- 延伸:Go 分布式事务实践:本地事务、Outbox、Saga、TCC 与 DTM
- 延伸:Go OpenTelemetry 实战:Trace、Metric、Context 与 OTLP
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论