如何使用paho mqtt客户端连接Digitransit指定主题MQTT broker
问题说明
需要使用Python的paho-mqtt客户端连接Digitransit MQTT broker获取公共交通实时数据,运行环境为浏览器端Jupyter Notebook,已知连接参数:
- Broker地址:
mqtt.hsl.fi - 端口:
8883 - 订阅主题:
/hfp/v2/journey/#
命令行环境下连接正常,可通过以下命令测试连通性:
- 安装命令行MQTT客户端
npm install -g mqtt
- 执行订阅命令
mqtt subscribe -h mqtt.hsl.fi -p 8883 -l mqtts -v -t "/hfp/v2/journey/#"
命令行可正常接收到如下格式的数据:
/hfp/v2/journey/ongoing/vp/bus/0022/01281/1040/2/Elielinaukio/14:29/1140118/0//// {"VP":{"desi":"40","dir":"2","oper":22,"veh":1281,"tst":"2022-07-06T11:56:53.417Z","tsi":1657108613,"spd":null,"hdg":null,"lat":null,"long":null,"acc":null,"dl":-101,"odo":8530,"drst":0,"oday":"2022-07-06","jrn":2646,"line":62,"start":"14:29","loc":"ODO","stop":null,"route":"1040","occu":0}}
相同逻辑的paho-mqtt代码在PyCharm、VS Code等本地IDE中可正常接收数据,但在浏览器端Jupyter Notebook中运行时,仅能触发on_connect回调打印连接成功信息,无法收到任何推送消息,测试代码如下:
# pip install paho-mqtt import paho.mqtt.client as mqtt # 收到服务端连接响应时的回调 def on_connect(client, userdata, flags, rc): print("Connected with result code "+ str(rc)) # 连接成功后订阅主题,断线重连时会自动重新订阅 client.subscribe("/hfp/v2/journey/#") # 收到推送消息时的回调 def on_message(client, userdata, msg): print(msg.topic+" "+str(msg.payload)) client = mqtt.Client() client.on_connect = on_connect client.on_message = on_message client.tls_set() client.connect("mqtt.hsl.fi", 8883, 60) client.loop_forever()
故障原因
核心问题出在事件循环的调用方式上:loop_forever()是阻塞式调用,启动后会一直占用当前主线程,阻塞后续所有逻辑执行。本地IDE运行Python脚本时,进程本身就是独立运行的,阻塞循环不会影响MQTT回调的触发;但浏览器端Jupyter Notebook的单元格执行、内核交互依赖主线程调度,loop_forever()的阻塞会导致消息处理回调没有机会被调度执行,因此只能看到连接成功的打印,无法收到后续消息。
修复方案
将阻塞式的loop_forever()替换为非阻塞的后台事件循环启动方法loop_start(),修改后的可运行代码如下:
# 安装依赖:pip install paho-mqtt import paho.mqtt.client as mqtt def on_connect(client, userdata, flags, rc): print(f"连接成功,返回码:{rc}") client.subscribe("/hfp/v2/journey/#") def on_message(client, userdata, msg): print(f"主题:{msg.topic}") print(f"消息内容:{msg.payload.decode('utf-8')}\n") client = mqtt.Client() client.on_connect = on_connect client.on_message = on_message # 8883端口为MQTT over TLS端口,必须配置TLS client.tls_set() client.connect("mqtt.hsl.fi", 8883, keepalive=60) # 启动后台非阻塞事件循环,替代阻塞的loop_forever() client.loop_start() # 需要停止接收消息时,执行 client.loop_stop() 即可
补充排查点:如果替换循环方法后仍无法收到消息,先检查浏览器端Jupyter部署环境的出站防火墙规则,确认8883端口的TLS连接没有被拦截。
内容的提问来源于stack exchange,提问作者Daniel
相关产品推荐
相关产品推荐

