Go操作Kafka的关键在于客户端选型、RequiredAcks和advertised.listeners配置、资源释放时机:kafka-go适合单分区验证但不支持消费者组自动重平衡,sarama适合生产环境需显式Close(),RequiredAcks须设WaitForAll且replication-factor≥2,advertised.listeners须正确配置IP而非localhost。

Go 操作 Kafka 不是“写个 WriteMessages 就完事”的事——默认配置下,消息可能根本没进磁盘、消费者会卡死、本地连不上 broker、重启后 offset 乱跳。关键在三件事:客户端选型、RequiredAcks 和 advertised.listeners 配置、资源释放时机。
选 kafka-go 还是 sarama?看你要不要消费者组自动重平衡
kafka-go 简洁,但 Reader 默认不支持消费者组自动重平衡;sarama 重一点,ConsumerGroup 实现完整,适合需 at-least-once 语义的生产环境。
- 单分区、快速验证、本地调试:用
kafka-go,WriteMessages一行发,ReadMessage一行收 - 对接已有集群、要 SSL/SASL/精确 offset 控制/消费者组:选
sarama,尤其注意它不自动关连接,Close()必须显式调用 - 别在
init()里全局初始化sarama.SyncProducer——它不是线程安全的,高并发下直接panic
sarama 生产者必须设的三个配置
NewSyncProducer 听起来可靠,但成功返回 ≠ 消息已落盘。Broker 默认配置下,RequiredAcks = sarama.NoResponse,发完就返回,broker 甚至没写入页缓存。
-
config.Producer.Return.Successes = true:否则你根本不知道消息发到了哪个 partition 和 offset -
config.Producer.RequiredAcks = sarama.WaitForAll(即acks=-1),但前提是 topic 的replication-factor >= 2,否则会退化为只等 leader - 必须显式调用
producer.Close(),否则 goroutine 泄露 + TCP 连接堆积,broker 日志会出现Connection reset
kafka-go 消费者为什么卡住?不是 Kafka 问题,是 offset 提交策略不对
ReadMessage 默认阻塞,但卡住往往不是 Kafka 问题,而是 offset 提交策略不当。
立即学习“go语言免费学习笔记(深入)”;
- 每条都立刻
CommitMessages:吞吐暴跌,RTT 成瓶颈;且提交后 crash,消息会被重复消费 - 完全不提交、只靠
AutoCommit: true:consumer 重启后从上次自动保存位置开始,可能跳过未处理完的消息 - 推荐做法:批量处理 10–100 条后手动
CommitMessages,同时用context.WithTimeout包裹ReadMessage,防无限阻塞
本地连不上 localhost:9092?先看 advertised.listeners
Docker 或远程 client 连不上,90% 是 broker 的 advertised.listeners 配置错了。Kafka broker 默认监听 localhost,但 client 会按 advertised.listeners 里的地址去连——如果它还是 localhost,就会连错。
- 本地开发:把
advertised.listeners设成PLAINTEXT://127.0.0.1:9092 - Docker 部署:设成宿主机可访问的 IP 或域名,比如
PLAINTEXT://host.docker.internal:9092 - SSL/SASL 场景下,还要确认
ssl.endpoint.identification.algorithm是否设为none(测试环境)或正确域名(生产)
最常被忽略的是:advertised.listeners 错了,所有客户端配置再对也白搭;sarama 不调 Close(),跑几天后连接数爆满;kafka-go 消费者没加 context.WithTimeout,一卡就是几小时没人发现。



















