消费者收不到消息主因是sarama.Config默认值过保守:必须设Consumer.Return.Errors=true、Consumer.Offsets.Initial为OffsetOldest/Newest、Net.DialTimeout和ReadTimeout≥10秒;流式处理优先选kafka-go.Reader,其无状态、单partition、精准offset提交更可靠。

消费者收不到消息?先检查 sarama.Config 这三个默认值
绝大多数“连上了但没消费到”的问题,不是 Kafka 配置错了,而是 sarama.Config 默认太保守,静默掩盖了失败。
-
Consumer.Return.Errors必须设为true,否则 offset 越界、重平衡失败等错误完全不抛出,只跳过消息 -
Consumer.Offsets.Initial不能依赖默认值(0),新消费者组首次启动必须显式设成sarama.OffsetOldest或sarama.OffsetNewest;否则 Kafka 查不到已提交 offset,直接返回UnknownMemberId并退订 -
Net.DialTimeout和Net.ReadTimeout在 Docker 或云环境里建议至少设为10 * time.Second,否则网络抖动导致连接被丢弃且不重试
用 kafka-go.Reader 还是 sarama.ConsumerGroup?看你的消费模式
流式处理(如实时日志、订单流水)优先选 kafka-go.Reader;需要复杂 rebalance 控制或跨 topic 协调的场景才考虑 sarama.ConsumerGroup。
-
kafka-go.Reader是无状态的,每个实例默认只读一个 partition,天然避免 CPU / I/O 倾斜;ReadMessage返回直接带Offset和Partition,业务处理完就能精准提交位点 -
sarama.ConsumerGroup的Setup/Cleanup方法一旦含阻塞操作(如 DB 连接池初始化),整个 rebalance 就卡住,表现为消费停滞、lag 持续上涨 -
kafka-go.Reader关键配置别漏:MinBytes: 1(防低频消息延迟)、CommitInterval: 1 * time.Second(panic 时不至于丢 offset)、StartOffset: kafka.FirstOffset(明确起点)
消息体怎么序列化才不怕升级?加 version 字段 + 避开 time.Time
Go struct 直接 json.Marshal 发 Kafka,版本迭代时字段增删极易让下游 panic 或静默丢数据——Kafka 不管 schema,而 Go 解析强类型。
在 Golang 中使用 samber/hot 进行内存缓存,支持 LRU、LFU、TinyLFU、W‑TinyLFU、S3FIFO、ARC、TwoQueue、SIEVE、FIFO 等淘汰算法,提供 TTL、缓存加载器及分片功能。
- 消息体最外层必须带
v字段,例如{"v": "1.2", "data": {...}},消费者按v分支解析,老版本逻辑可 fallback - 绝对不要把
time.Time字段直接塞进 JSON:不同 Go 版本、不同语言消费者对零值、时区、格式的处理不一致;统一转成 RFC3339 字符串或 Unix 毫秒整数 - 用 Protobuf 的话,别手写序列化逻辑;两端共用同一份
.proto文件,并用gogoproto+gofast插件生成代码
生产者别用 sarama.SyncProducer,改用 sarama.AsyncProducer
sarama.SyncProducer 是阻塞式,超时难控制、不支持批量、失败后需手动重试,不适合微服务高频发事件场景。
立即学习“go语言免费学习笔记(深入)”;
- 启用
Producer.Idempotent = true,获得单 partition 内的幂等性(避免重复投递),这是目前 sarama v1.35 支持的最接近 Exactly-Once 的能力 - 监听
Successes()和Errors()两个 channel 处理结果,别只看SendMessage返回值 - topic 名要带 service 和 domain event 语义,比如
order-service.order-created.v1,而不是模糊的orders——方便后续权限隔离、监控追踪和 schema 管理
实际跑起来之后,最容易被忽略的是:消费者 offset 提交时机和 panic 恢复边界是否对齐。哪怕用了 kafka-go.Reader,如果业务 handler 里 panic 了又没 recover,CommitInterval 机制就失效;而 sarama 的 session.Commit() 更容易被中断。这事关数据是否丢失,不是配完就完事。


















