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

K8S环境下MQTT转GCP IoT Core中继程序运行异常排查

问题根源分析

你的程序在K8S中运行一段时间后停止处理消息,核心问题出在relay函数的逻辑设计上,导致goroutine阻塞和通道死锁:

关键问题点

在relay函数中,当你从jwtTokens通道拿到第一个JWT令牌后,会进入一个无限内部循环持续读取messages通道的消息并转发:

for token := range jwtTokens {
    log.Println("Connecting with fresh token")
    client := googleConnect(clientId, googleUri, token)
    // 这里的无限循环没有退出条件
    for {
        msg := <-messages
        // 消息转发逻辑...
    }
}

这个内部循环会永久占用当前goroutine,导致外层的for token := range jwtTokens无法处理后续生成的新JWT令牌。当messages通道被填满(无缓冲或缓冲耗尽)时,负责接收MQTT消息的listen goroutine会因为往messages通道写入数据被阻塞,最终整个程序陷入停滞。

同时,每次生成新JWT后,旧的Google IoT Core客户端连接没有被主动断开,还会造成连接资源泄漏。

修复方案

我们需要给内部消息处理循环添加退出信号,当新的JWT令牌到来时,优雅地终止旧循环、断开旧连接,再用新令牌建立新连接继续处理消息。

修改后的relay函数代码如下:

// relay Relays the messages handled by listen to Google IoT Core
// Connection is refreshed whenever the JWT is refreshed
func relay(messages chan mqtt.Message, jwtTokens chan string) {
    googleUri, err := url.Parse("mqtt://mqtt.googleapis.com:443/")
    if err != nil {
        log.Fatal(err)
    }
    clientIdSlice := []string{"projects", "domalys-202111", "locations", "europe-west1", "registries", "domalys-mqtt", "devices", "mqtt-relay"}
    clientId := strings.Join(clientIdSlice, "/")

    var quitChan chan struct{} // 用于通知旧循环退出
    var currentClient mqtt.Client

    for token := range jwtTokens {
        log.Println("Connecting with fresh token")
        // 如果存在旧连接,先断开并发送退出信号
        if quitChan != nil {
            close(quitChan)
            if currentClient != nil && currentClient.IsConnected() {
                currentClient.Disconnect(250)
                log.Println("Disconnected old Google IoT Core client")
            }
        }

        // 创建新的退出通道和客户端
        quitChan = make(chan struct{})
        currentClient = googleConnect(clientId, googleUri, token)

        // 启动新的消息处理goroutine
        go func(client mqtt.Client, quit <-chan struct{}) {
            for {
                select {
                case msg := <-messages:
                    payload := string(msg.Payload())
                    topic := "/devices/mqtt-relay/events/" + msg.Topic()
                    log.Println("Sending : ", payload)
                    client.Publish(topic, 1, false, payload)
                case <-quit:
                    log.Println("Stopping old message processing loop")
                    return
                }
            }
        }(currentClient, quitChan)
    }
}

修复说明

  1. 退出通道机制:每次新JWT令牌到来时,关闭旧的quitChan,通知正在运行的消息处理goroutine退出。
  2. 独立goroutine处理消息:将消息转发逻辑放到独立的goroutine中,这样外层的for token := range jwtTokens可以持续处理新的JWT令牌。
  3. 资源清理:在切换新连接前,主动断开旧的Google IoT Core客户端连接,避免资源泄漏。

额外优化建议

  • 在listen函数中添加重连逻辑:如果本地MQTT服务器连接断开,当前代码会直接退出程序,你可以在connect函数中添加循环重连的逻辑,提升程序鲁棒性。
  • 给createJWT函数的错误处理补全:当前generateJwtToken函数中如果生成JWT失败,只会忽略错误,建议添加错误日志,方便排查问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 21:18:15