You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用segmentio/kafka-go编写Golang Kafka消费者出现i/o timeout超时错误

问题根因
  • 核心报错指向Kafka服务端广播的可访问地址异常:你代码中配置的连接地址是localhost:9092,但消费者首次连接Broker后,会从Broker拿到其配置的对外广播地址advertised.listeners用于后续拉取消息,你的服务端该配置默认对应了127.0.1.1:9092,本地网络无法正常访问该地址导致超时
  • 消费者配置的MaxBytes参数值过低:10字节的上限远小于常规消息体积,Broker无法返回符合大小要求的消息,也会触发读取超时
  • 主进程无阻塞逻辑:main函数启动消费者协程后没有阻塞逻辑,主进程可能提前退出终止消费流程
修复步骤
  1. 修正Kafka服务端配置
    找到Kafka安装目录下config/server.properties文件,修改以下两个配置项:
# 配置服务端监听所有网卡的9092端口
listeners=PLAINTEXT://0.0.0.0:9092
# 配置对外广播的访问地址,消费者最终会用这个地址和Broker交互
advertised.listeners=PLAINTEXT://localhost:9092

修改完成后重启Kafka服务生效。

  1. 调整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 {} // 永久阻塞主进程,避免消费协程被终止
}
  1. 验证服务可用性
    修改完成后先通过Kafka自带CLI工具验证消费能力,确认消息可正常读取:
# 替换为你环境中对应脚本的路径和执行方式
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic firsttopic --from-beginning

CLI验证通过后再运行Go代码即可正常消费消息。

内容的提问来源于stack exchange,提问作者Vivek Jaiswal

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.07 09:24:04