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

如何使用paho mqtt客户端连接Digitransit指定主题MQTT broker

问题说明

需要使用Python的paho-mqtt客户端连接Digitransit MQTT broker获取公共交通实时数据,运行环境为浏览器端Jupyter Notebook,已知连接参数:

  • Broker地址:mqtt.hsl.fi
  • 端口:8883
  • 订阅主题:/hfp/v2/journey/#

命令行环境下连接正常,可通过以下命令测试连通性:

  1. 安装命令行MQTT客户端
npm install -g mqtt
  1. 执行订阅命令
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 14:09:15