绝大多数“收不到消息”是因消费者组配置未对齐:group.id非法、metadata重试不足、offset起始位置设为newest导致跳过历史消息、分区分配为空、未手动调用markmessage提交offset、本地docker环境bootstrap.servers误用localhost等。

sarama.NewConsumerGroup 为什么收不到消息
绝大多数“收不到消息”不是代码写错,而是消费者组初始化或运行时配置没对齐。Kafka 只把新消息推给有合法 offset 位置的 group 成员,不自动回溯历史。
- group.id 为空或含非法字符(如空格、开头下划线、超 249 字节)→ Kafka 拒绝注册,日志只报
group coordinator not available - topic 不存在且
config.Metadata.Retry.Max太小(默认 3 次)→ 首次 fetch metadata 失败就退出,建议显式设为5 - 启动时用
sarama.OffsetNewest,但消息在 consumer 启动前已发出 → 它只收之后的新消息,看起来像“一直没数据” - topic 分区数 > 1,但当前 consumer 实例只被分配到空分区 → 查看
ConsumeClaim中claim.Partition()和claim.HighWaterMarkOffset()确认是否真有数据可读
ConsumerGroupHandler 的 ConsumeClaim 必须手动 MarkMessage
session.MarkMessage(message, "") 不是可选项,是 offset 提交的前提。不调它,哪怕消息处理完了,offset 也不会提交,下次重启照样重拉一遍。
- 必须在
ConsumeClaim函数体内、消息处理成功后立即调用,不能丢到 goroutine 里异步执行 - 如果处理失败需跳过,也得调
session.MarkMessage(message, "")或session.MarkOffset(claim.Topic(), claim.Partition(), message.Offset+1, ""),否则会卡住整个分区 - 别依赖
config.Consumer.Offsets.AutoCommit.Enable = true:默认 1s 提交一次,延迟高、不可控;生产环境建议关掉,自己控制提交时机
本地开发连 Docker Kafka 时 bootstrap.servers 怎么填
写 localhost:9092 是最常见错误。Go 进程若跑在容器里,localhost 指的是容器自身,不是宿主机上的 Kafka。
- Docker Desktop / Colima 环境:改用
host.docker.internal:9092(Mac/Win 支持,Linux 需额外配置) - Linux Docker:用宿主机真实 IP,如
192.168.1.100:9092,并确认 Kafka 的advertised.listeners已配成该地址 - 若 Kafka 也在容器中(如 docker-compose),则填服务名 + 端口,如
kafka:9092,并确保网络互通 - 无论哪种,启动前先用
nc -zv $HOST $PORT测试连通性,别等跑起来再猜
ConsumerGroup 初始化失败的三个硬伤点
调 sarama.NewConsumerGroup 返回 nil 或 panic,基本不是 Go 代码逻辑问题,而是底层连接/认证/协议层面卡住了。
- broker 地址解析失败:Go 客户端无法解析容器名或 DNS 别名 → 改用 IP + 显式端口,或检查容器网络模式(
host模式更简单) - SASL/SSL 认证缺失:集群开了 SASL_PLAIN 而 config 没配
config.Net.SASL.Enable = true→ 日志出现failed to find SASL mechanism - Kafka 协议版本不匹配:比如集群是 3.6.0,却设了
config.Version = sarama.V3_7_0_0(sarama 当前最高只支持到V3_6_0_0)→ 握手失败,静默断连或报invalid request type
session.Context() 的生命周期和 claim.Messages() 的阻塞行为——它们共同决定了 consumer 是否能及时响应 rebalance 或优雅退出。没处理好,就会出现“明明程序停了,Kafka 还以为它活着”,导致分区长时间无法再均衡。golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











