一句话总结
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.理由:
- 默认幂等 + acks=all,不用操心可靠性配置
- API 设计贴合 Go 习惯(context、error、functional options)
- 纯 Go,CI/CD 和容器化零摩擦
- 维护者响应极快,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"), ) 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"),
kgo.RequiredAcks(kgo.AllISRAcks()),
kgo.ProducerBatchMaxBytes(256 * 1024), kgo.ProducerLinger(10 * time.Millisecond),
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), ) 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"),
kgo.ConsumerGroup("order-processor"), kgo.ConsumeTopics("order-events"),
kgo.DisableAutoCommit(),
kgo.ConsumeResetOffset(kgo.NewOffset().AtEnd()),
kgo.Balancers(kgo.CooperativeStickyBalancer()),
kgo.SessionTimeout(45 * time.Second),
kgo.BlockRebalanceOnPoll(), )
|
消费循环
func consumeLoop(ctx context.Context, client *kgo.Client) error { for { fetches := client.PollFetches(ctx)
if ctx.Err() != nil { return ctx.Err() }
if errs := fetches.Errors(); len(errs) > 0 { for _, e := range errs { slog.Error("fetch error", "topic", e.Topic, "partition", e.Partition, "err", e.Err, ) } }
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, ) } } }(p.Records) }) wg.Wait()
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——你需要:
- 停止拉取新消息
- 等待正在处理的消息处理完
- 提交最后的 offset
- 通知 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...")
cancel()
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 里做重操作
client.Produce(ctx, record, func(r *kgo.Record, err error) { time.Sleep(5 * time.Second) db.Insert(r) })
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): produceToDLQ(r, err) return nil
default: slog.Error("unknown error, skipping", "err", err) return nil } }
|
3. 健康检查
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 的工程实现.