Kafka 不是 Go 微服务默认消息总线,需显式处理序列化、消费者组管理、错误重试和 offset 提交;生产环境优先选 kafka-go(支持 context.Context、优雅停机),但须手动提交 offset;发消息必须带重试+退避。

Kafka 不是 Go 微服务的默认消息总线,接入它必须显式处理序列化、消费者组管理、错误重试和 offset 提交策略——跳过这些环节,服务上线后大概率出现消息丢失或重复消费。
用 sarama 还是 kafka-go?选错库会卡在 reconnect 或 context cancel 上
生产环境优先选 kafka-go(segmentio 出品):它原生支持 context.Context,消费者能响应 ctx.Done() 快速退出;sarama 的 ConsumePartition 是阻塞调用,没封装好容易导致服务无法优雅停机。
但注意:kafka-go 默认不自动提交 offset,必须手动调用 conn.CommitOffsets 或启用 AutoCommit 配置;而 sarama 的 ConsumerGroup 实现更贴近 Kafka 原语,适合需要精细控制 rebalance 行为的场景。
- 新项目、追求简洁和 context 友好 → 用
kafka-go,配kafka.NewReader+ReadMessage - 已有 sarama 使用经验、需自定义 heartbeat 或 metadata 请求逻辑 → 继续用
sarama,但务必封装ConsumerGroupHandler的Setup/Teardown方法 - 别直接用
sarama.AsyncProducer发送消息:它不保证发送成功,错误只走Errors()channel,容易漏处理
kafka-go 消费者必须自己处理 offset 提交时机,否则重启就丢消息
默认配置下,kafka-go 的 Reader 不自动提交 offset。如果只调 reader.ReadMessage(ctx) 就结束,下次启动会从上次提交的位置继续读——而那个位置可能是几小时前的。
立即学习“go语言免费学习笔记(深入)”;
正确做法是:在业务逻辑执行成功后,立即提交当前 message 的 offset:
Colly 是一个用于 Go 语言的快速开源爬取和爬虫框架。它适用于从简单的页面提取到异步爬虫处理大量页面集合,支持请求回调和结构化解析。
msg, err := reader.ReadMessage(ctx)
if err != nil {
return err
}
// 处理业务逻辑
if err := process(msg.Value); err != nil {
return err
}
// ✅ 成功后提交 offset
if err := reader.CommitMessages(ctx, msg); err != nil {
log.Printf("commit failed: %v", err)
return err
}
- 不要在 goroutine 里异步提交 offset:可能提交了还没处理完,或者处理失败了却已提交
- 避免批量提交(如 accumulate 10 条再 commit):增加重复消费概率,且故障时回溯困难
- 如果用
reader.Config().MaxWait调大等待时间,要同步调高session.timeout.ms对应的客户端参数,否则触发 rebalance
Go 微服务发消息必须带重试+退避,Kafka broker 临时不可用很常见
Kafka 写入失败不是异常,而是常态:网络抖动、broker 重启、磁盘满都会返回 NetworkException 或 NotEnoughReplicasException。裸调 kafka.Writer.WriteMessages 不做重试,等于把可靠性交给运气。
推荐组合:kafka.Writer + 自定义重试逻辑(不用第三方 retry 库,避免 context 泄漏):
for i := 0; i < 3; i++ {
err := writer.WriteMessages(ctx, kafka.Message{Value: data})
if err == nil {
return nil
}
if !isRetriable(err) {
return err
}
time.Sleep(time.Second * time.Duration(1<<uint(i))) // 1s, 2s, 4s
}
return fmt.Errorf("failed after 3 retries: %w", err)
-
isRetriable至少覆盖*kafka.UnknownTopicOrPartitionError、*kafka.NetworkException、*kafka.NotEnoughReplicasException - 别用
time.AfterFunc做重试:goroutine 生命周期难管理,容易堆积 - 写入前检查
ctx.Err(),防止重试过程中服务已关闭
微服务间消息格式必须约定 schema,别传裸 JSON 或 map[string]interface{}
看似方便的 json.Marshal(map[string]interface{}) 会导致下游解析失败:字段类型不一致(比如 int vs string)、字段缺失无提示、新增字段无法向后兼容。
真实项目中,应该用 avro 或至少固定结构的 Go struct + json 标签:
type OrderCreatedEvent struct {
OrderID string `json:"order_id"`
UserID int64 `json:"user_id"`
Total float64 `json:"total"`
CreatedAt time.Time `json:"created_at"`
}
- 所有事件 struct 必须加
json:tag,避免大小写不一致导致反序列化为零值 - 时间字段统一用
time.Time,别用字符串或 Unix 时间戳:时区和精度问题会在多个服务间放大 - 如果用 Avro,schema registry 地址必须作为配置项注入,不能硬编码在代码里
最常被忽略的一点:Kafka 的 message.Key 不只是用来分区,它决定了同一 key 的消息在 topic 内严格有序——如果你的订单状态更新依赖顺序,但发消息时没设 Key,那下游看到的状态流转就是乱序的。


















