阶段一 · 基础与核心模型

Go 客户端实战:sarama,confluent-kafka-go,franz-go

一句话总结

Go 生态有三个主流 Kafka 客户端:sarama(最老,社区最大)、confluent-kafka-go(C 绑定,性能最强)、franz-go(纯 Go,API 最现代).
本篇对比三者取舍,给出 franz-go 的完整工程结构:Producer 服务、Consumer 服务、优雅关停、错误处理.

三大客户端对比

维度 sarama confluent-kafka-go franz-go
实现方式 纯 Go cgo 封装 librdkafka(C 库) 纯 Go
维护方 IBM/Shopify → 社区 Confluent 官方 Travis Bischel(个人,极活跃)
API 风格 低层级,需手动管理大量细节 中层级,偏 C 风格回调 高层级,Go 惯用法
幂等/事务 需手动配置,坑多 完整支持 默认开启幂等
交叉编译 ✓(纯 Go) ✗(依赖 C 库,CGO_ENABLED=1) ✓(纯 Go)
性能 中等 最高(C 库优化) 接近 confluent,远超 sarama
Consumer Group 手动实现 ConsumerGroupHandler 接口 回调风格 一个 Client 搞定,极简
文档质量 中等,靠社区博客 好,官方文档 优秀,源码注释极详细

选型建议

场景 推荐 理由
新项目,追求开发效率 franz-go API 最简洁,默认配置最安全,纯 Go 无 cgo 烦恼
已有 sarama 代码在跑 继续 sarama 或渐进迁移 迁移成本高于收益,除非遇到 sarama 的 bug
极致吞吐 + 运维团队能搞定 cgo confluent-kafka-go librdkafka 久经沙场,吞吐天花板最高
需要交叉编译(如 ARM/Alpine) franz-go 或 sarama 纯 Go,GOOS/GOARCH 直接 build

本系列选 franz-go

后续所有代码示例使用 franz-go.理由:

  1. 默认幂等 + acks=all,不用操心可靠性配置
  2. API 设计贴合 Go 习惯(context、error、functional options)
  3. 纯 Go,CI/CD 和容器化零摩擦
  4. 维护者响应极快,issue 通常当天回复

安装与初始化

go get github.com/twmb/franz-go/pkg/kgo

franz-go 的核心是一个 kgo.Client,同时具备 Producer 和 Consumer 能力——不像 sarama 要分别创建 SyncProducer 和 ConsumerGroup.

import "github.com/twmb/franz-go/pkg/kgo"

client, err := kgo.NewClient(
kgo.SeedBrokers("broker1:9092", "broker2:9092", "broker3:9092"),
// 更多 options 按需添加
)
if err != nil {
log.Fatal("create kafka client", "err", err)
}
defer client.Close()

kgo.NewClient 接受一组 functional options,下面按场景分别展开.


Producer 完整示例

基础配置

client, err := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.DefaultProduceTopic("order-events"),

// 可靠性(franz-go 默认已开启幂等 + acks=all,以下是显式声明)
kgo.RequiredAcks(kgo.AllISRAcks()),

// 批量优化
kgo.ProducerBatchMaxBytes(256 * 1024), // 单批上限 256KB
kgo.ProducerLinger(10 * time.Millisecond), // 凑批等待 10ms

// 压缩
kgo.ProducerBatchCompression(kgo.Lz4Compression()),
)

异步发送(主流用法)

func produceOrderEvent(ctx context.Context, client *kgo.Client, userID string, payload []byte) {
record := &kgo.Record{
Key: []byte(userID),
Value: payload,
}

client.Produce(ctx, record, func(r *kgo.Record, err error) {
if err != nil {
// 所有内部重试都失败了才到这里
slog.Error("produce failed",
"err", err,
"topic", r.Topic,
"key", string(r.Key),
)
// 业务决策:写死信队列 / 报警 / metric++
return
}
slog.Debug("produced",
"topic", r.Topic,
"partition", r.Partition,
"offset", r.Offset,
)
})
}

同步发送(低频关键事件)

results := client.ProduceSync(ctx, record)
if err := results.FirstErr(); err != nil {
return fmt.Errorf("produce sync: %w", err)
}

Consumer 完整示例

基础配置

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

// Consumer Group 配置
kgo.ConsumerGroup("order-processor"),
kgo.ConsumeTopics("order-events"),

// 关闭自动提交,手动控制
kgo.DisableAutoCommit(),

// 从最新消息开始(新 Group 首次消费时)
kgo.ConsumeResetOffset(kgo.NewOffset().AtEnd()),

// Cooperative rebalance(默认,显式声明便于理解)
kgo.Balancers(kgo.CooperativeStickyBalancer()),

// 防止处理太慢被踢
kgo.SessionTimeout(45 * time.Second),

// 阻止 rebalance 打断处理中的消息
kgo.BlockRebalanceOnPoll(),
)

消费循环

func consumeLoop(ctx context.Context, client *kgo.Client) error {
for {
fetches := client.PollFetches(ctx)

// context 取消时退出
if ctx.Err() != nil {
return ctx.Err()
}

// 处理 fetch 错误(网络问题、鉴权失败等)
if errs := fetches.Errors(); len(errs) > 0 {
for _, e := range errs {
slog.Error("fetch error",
"topic", e.Topic,
"partition", e.Partition,
"err", e.Err,
)
}
// 大部分 fetch error 是暂时的,继续下一轮 poll
}

// 按 Partition 并行处理(保证分区内有序)
var wg sync.WaitGroup
fetches.EachPartition(func(p kgo.FetchTopicPartition) {
wg.Add(1)
go func(records []*kgo.Record) {
defer wg.Done()
for _, r := range records {
if err := processRecord(r); err != nil {
slog.Error("process failed",
"err", err,
"topic", r.Topic,
"partition", r.Partition,
"offset", r.Offset,
)
// 策略:跳过 or 重试 or 写死信
}
}
}(p.Records)
})
wg.Wait()

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

func processRecord(r *kgo.Record) error {
// 你的业务逻辑
slog.Info("processing",
"key", string(r.Key),
"partition", r.Partition,
"offset", r.Offset,
)
return nil
}

优雅关停(Graceful Shutdown)

Kafka Consumer 的关停不能直接 os.Exit——你需要:

  1. 停止拉取新消息
  2. 等待正在处理的消息处理完
  3. 提交最后的 offset
  4. 通知 Group Coordinator 主动离组(避免等 session timeout 才触发 rebalance)
func main() {
ctx, cancel := context.WithCancel(context.Background())

client, err := kgo.NewClient(/* ... */)
if err != nil {
log.Fatal(err)
}

// 启动消费
go func() {
if err := consumeLoop(ctx, client); err != nil && err != context.Canceled {
slog.Error("consume loop exited", "err", err)
}
}()

// 等待退出信号
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
<-sigCh

slog.Info("shutting down...")

// 1. 取消 context → consumeLoop 的 PollFetches 返回
cancel()

// 2. Close 会:
// - 等待 in-flight produce 回调完成
// - 提交 uncommitted offsets(如果配了 autocommit)
// - 向 Coordinator 发送 LeaveGroup
client.Close()

slog.Info("shutdown complete")
}

client.Close() 做了什么

Close() 内部会:

  • 刷空 Producer 缓冲区(等待 ack)
  • 如果有 Consumer Group,发送 LeaveGroup 请求
  • 关闭所有 TCP 连接

所以 defer client.Close() 是必须的——不调用的话,Group 要等 session.timeout.ms(默认 45s)才会发现你离开,期间该 Consumer 的 Partition 没人消费.


完整工程结构

一个典型的 Kafka Consumer 微服务目录:

order-consumer/
├── main.go // 启动入口,信号处理,DI 组装
├── internal/
│ ├── kafka/
│ │ ├── client.go // NewClient 封装,配置集中管理
│ │ ├── consumer.go // consumeLoop,按 partition 分发
│ │ └── producer.go // produce 封装(如果本服务也需要生产)
│ ├── handler/
│ │ └── order.go // 业务处理逻辑(processRecord 的实现)
│ └── config/
│ └── config.go // 配置结构体(从环境变量/配置中心加载)
├── go.mod
└── go.sum

配置管理

把 Kafka 配置抽成结构体,方便从环境变量或配置中心加载:

type KafkaConfig struct {
Brokers []string `yaml:"brokers"`
Topic string `yaml:"topic"`
GroupID string `yaml:"group_id"`
SessionTimeout time.Duration `yaml:"session_timeout"`
LingerMs time.Duration `yaml:"linger_ms"`
BatchMaxBytes int32 `yaml:"batch_max_bytes"`
}

func NewClient(cfg KafkaConfig) (*kgo.Client, error) {
opts := []kgo.Opt{
kgo.SeedBrokers(cfg.Brokers...),
kgo.ConsumerGroup(cfg.GroupID),
kgo.ConsumeTopics(cfg.Topic),
kgo.DisableAutoCommit(),
kgo.BlockRebalanceOnPoll(),
kgo.SessionTimeout(cfg.SessionTimeout),
kgo.Balancers(kgo.CooperativeStickyBalancer()),
}

if cfg.LingerMs > 0 {
opts = append(opts, kgo.ProducerLinger(cfg.LingerMs))
}
if cfg.BatchMaxBytes > 0 {
opts = append(opts, kgo.ProducerBatchMaxBytes(cfg.BatchMaxBytes))
}

return kgo.NewClient(opts...)
}

常见陷阱与最佳实践

1. 不要在 callback 里做重操作

// ✗ 错误:callback 在 franz-go 内部 goroutine 执行,阻塞会影响后续 batch 发送
client.Produce(ctx, record, func(r *kgo.Record, err error) {
time.Sleep(5 * time.Second) // 阻塞内部 sender
db.Insert(r) // 又一个阻塞操作
})

// ✓ 正确:callback 只做轻量操作(记日志、metric、塞 channel)
client.Produce(ctx, record, func(r *kgo.Record, err error) {
if err != nil {
metrics.ProduceErrors.Inc()
errCh <- ProduceError{Record: r, Err: err}
return
}
metrics.ProduceSuccess.Inc()
})

2. Consumer 处理失败的策略

func processRecord(r *kgo.Record) error {
err := businessLogic(r)
if err == nil {
return nil
}

// 策略选择:
switch {
case errors.Is(err, ErrTransient):
// 暂时性错误:本地重试几次
return retry(3, func() error { return businessLogic(r) })

case errors.Is(err, ErrPermanent):
// 永久性错误:写入死信 Topic,跳过这条
produceToDLQ(r, err)
return nil

default:
// 未知错误:记录并跳过,别卡住整个消费
slog.Error("unknown error, skipping", "err", err)
return nil
}
}

3. 健康检查

// 用 Ping 验证连接可用性(适合 K8s liveness/readiness probe)
func healthCheck(client *kgo.Client) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), 3*time.Second)
defer cancel()

if err := client.Ping(ctx); err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
fmt.Fprintf(w, "kafka unhealthy: %v", err)
return
}
w.WriteHeader(http.StatusOK)
w.Write([]byte("ok"))
}
}

从 sarama 迁移的关键差异

如果你的项目正在用 sarama,以下是最大的心智模型差异:

维度 sarama franz-go
Client 角色 Producer 和 Consumer 是不同对象 一个 Client 同时具备两种能力
Consumer Group 实现 ConsumerGroupHandler 三个方法 直接 PollFetches() 循环
Offset 提交 MarkOffset + CommitOffsets CommitUncommittedOffsets(自动追踪已 poll 的 offset)
Rebalance 感知 Setup/Cleanup 回调 OnPartitionsAssigned/Revoked option 或 BlockRebalanceOnPoll
幂等 手动设置 Producer.Idempotent = true + 一堆关联配置 默认开启,零配置
错误处理 Errors() channel(容易漏读导致死锁) callback 参数 / PollFetches().Errors()

小结

要点 说明
选 franz-go 纯 Go、默认安全配置、API 简洁、维护活跃
一个 Client 搞定 不需要分别创建 Producer/Consumer 对象
异步 Produce + callback 高吞吐主流用法,callback 里只做轻量操作
PollFetches + 手动 commit Consumer 黄金模式,配合 BlockRebalanceOnPoll
优雅关停 cancel context → client.Close()(刷缓冲 + LeaveGroup)
工程化 配置抽结构体、错误分级处理、健康检查暴露给 K8s

下一篇讨论消息可靠性的完整图景:端到端的不丢、不重、顺序性保证,以及 At-least-once 与 Exactly-once 的工程实现.