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) } }
修复说明
- 退出通道机制:每次新JWT令牌到来时,关闭旧的
quitChan,通知正在运行的消息处理goroutine退出。 - 独立goroutine处理消息:将消息转发逻辑放到独立的goroutine中,这样外层的
for token := range jwtTokens可以持续处理新的JWT令牌。 - 资源清理:在切换新连接前,主动断开旧的Google IoT Core客户端连接,避免资源泄漏。
额外优化建议
- 在
listen函数中添加重连逻辑:如果本地MQTT服务器连接断开,当前代码会直接退出程序,你可以在connect函数中添加循环重连的逻辑,提升程序鲁棒性。 - 给
createJWT函数的错误处理补全:当前generateJwtToken函数中如果生成JWT失败,只会忽略错误,建议添加错误日志,方便排查问题。
内容的提问来源于stack exchange,提问作者MoskitoHero
相关产品推荐
相关产品推荐

