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

如何用Golang实现Paho MQTT永久监听(类似Python的loop_forever)

在Golang中实现MQTT永久监听消息

在Golang常用的MQTT客户端库github.com/eclipse/paho.mqtt.golang中,没有Python Paho里loop_forever()那样的直接方法,但可以通过以下几种方式实现永久监听消息的效果:

方法一:用client.Wait()阻塞主线程

客户端连接成功后,调用client.Wait()会一直阻塞,直到客户端断开连接,这是最简单的实现方式:

package main

import (
    "fmt"
    "os"

    mqtt "github.com/eclipse/paho.mqtt.golang"
)

// 消息处理函数
var messageHandler mqtt.MessageHandler = func(client mqtt.Client, msg mqtt.Message) {
    fmt.Printf("收到消息: 主题=%s, 内容=%s\n", msg.Topic(), string(msg.Payload()))
}

// 连接成功回调
var onConnectHandler mqtt.OnConnectHandler = func(client mqtt.Client) {
    fmt.Println("已连接到MQTT服务器")
    // 订阅目标主题
    token := client.Subscribe("your/topic", 1, messageHandler)
    if token.Wait() && token.Error() != nil {
        fmt.Println("订阅失败:", token.Error())
        os.Exit(1)
    }
}

// 连接断开回调
var onConnLostHandler mqtt.ConnectionLostHandler = func(client mqtt.Client, err error) {
    fmt.Printf("连接断开: %v\n", err)
}

func main() {
    // 配置客户端选项
    opts := mqtt.NewClientOptions().AddBroker("tcp://mqtt.example.com:1883")
    opts.SetClientID("golang-mqtt-listener")
    opts.SetDefaultPublishHandler(messageHandler)
    opts.OnConnect = onConnectHandler
    opts.OnConnectionLost = onConnLostHandler

    // 创建并连接客户端
    client := mqtt.NewClient(opts)
    token := client.Connect()
    if token.Wait() && token.Error() != nil {
        panic(token.Error())
    }

    // 阻塞主线程,保持监听
    client.Wait()
}

方法二:结合信号监听实现优雅退出

如果需要支持优雅关闭(比如捕获Ctrl+C信号),可以通过监听系统信号配合客户端的Done()通道实现:

package main

import (
    "fmt"
    "os"
    "syscall"
    "os/signal"

    mqtt "github.com/eclipse/paho.mqtt.golang"
)

// 消息处理、连接回调函数同方法一...

func main() {
    // ... 客户端配置和连接代码同方法一 ...

    // 监听中断信号
    sigChan := make(chan os.Signal, 1)
    signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)

    // 阻塞等待信号或客户端断开
    select {
    case <-sigChan:
        fmt.Println("收到退出信号,正在断开连接...")
        client.Disconnect(250) // 等待250ms处理剩余消息
    case <-client.Done():
        fmt.Println("MQTT客户端已断开")
    }
}

核心要点

  • 客户端内部会自动启动后台goroutine处理消息收发,只要连接保持,就会持续监听
  • client.Wait()本质是等待client.Done()通道关闭,客户端断开或调用Disconnect()时该通道会关闭
  • 订阅操作放在OnConnect回调中,能确保连接重连后自动重新订阅主题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 23:18:34