使用segmentio/kafka-go编写Golang Kafka消费者出现i/o timeout超时错误
问题根因
- 核心报错指向Kafka服务端广播的可访问地址异常:你代码中配置的连接地址是
localhost:9092,但消费者首次连接Broker后,会从Broker拿到其配置的对外广播地址advertised.listeners用于后续拉取消息,你的服务端该配置默认对应了127.0.1.1:9092,本地网络无法正常访问该地址导致超时 - 消费者配置的
MaxBytes参数值过低:10字节的上限远小于常规消息体积,Broker无法返回符合大小要求的消息,也会触发读取超时 - 主进程无阻塞逻辑:main函数启动消费者协程后没有阻塞逻辑,主进程可能提前退出终止消费流程
修复步骤
- 修正Kafka服务端配置
找到Kafka安装目录下config/server.properties文件,修改以下两个配置项:
# 配置服务端监听所有网卡的9092端口 listeners=PLAINTEXT://0.0.0.0:9092 # 配置对外广播的访问地址,消费者最终会用这个地址和Broker交互 advertised.listeners=PLAINTEXT://localhost:9092
修改完成后重启Kafka服务生效。
- 调整Go消费者配置
修正MaxBytes参数,可选增加超时、偏移量配置,调整后的代码参考:
import "time" // 需要导入time包 func Startkafka() { conf := kafka.ReaderConfig{ Brokers: []string{"localhost:9092"}, Topic: "firsttopic", GroupID: "g1", MaxBytes: 10240, // 调整为10KB,根据实际消息大小调整即可 DialTimeout: 10 * time.Second, // 连接超时配置 ReadTimeout: 10 * time.Second, // 读取超时配置 StartOffset: kafka.FirstOffset, // 新消费组默认从头消费历史消息,可按需删除 } reader := kafka.NewReader(conf) for { m, err := reader.ReadMessage(context.Background()) if err != nil { fmt.Println("Some error occured", err) continue } fmt.Println("Message is : ", string(m.Value)) } } func main() { go Startkafka() fmt.Println("Kafka has been started...") select {} // 永久阻塞主进程,避免消费协程被终止 }
- 验证服务可用性
修改完成后先通过Kafka自带CLI工具验证消费能力,确认消息可正常读取:
# 替换为你环境中对应脚本的路径和执行方式 kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic firsttopic --from-beginning
CLI验证通过后再运行Go代码即可正常消费消息。
内容的提问来源于stack exchange,提问作者Vivek Jaiswal
相关产品推荐
相关产品推荐

