如何实现常规MQTT与Tornado WebSocket通信?代码问题求助
问题:MQTT消息无法通过Tornado WebSocket转发到浏览器
现有代码中Tornado服务运行正常,单独测试MQTT监听功能也能正常接收消息,但MQTT收到的消息无法传递给已连接的WebSocket客户端,目标是实现将MQTT消息转发至浏览器端。
问题原因
Paho-MQTT的loop_start()会启动独立后台线程处理MQTT消息回调,而Tornado的WebSocket操作(如write_message)必须在Tornado自身的IO循环线程中执行,跨线程直接调用WebSocket方法会导致消息无法正确投递。
解决方案
修改消息发送逻辑,通过Tornado的IO循环add_callback()方法,将WebSocket消息发送操作提交到正确的线程中执行。
修改后的完整代码
import os.path import tornado.httpserver import tornado.web import tornado.ioloop import tornado.options import tornado.httpclient import tornado.websocket import json import random from paho.mqtt import client as mqtt_client broker = 'xxxxxxx' port = 1883 topic = "xxx/xxxx" client_id = f'xxxx-{random.randint(0, 100)}' clients = [] def connect_mqtt() -> mqtt_client: def on_connect(clientMQTT, userdata, flags, rc): if rc == 0: print("Connected to MQTT Broker!") else: print("Failed to connect, return code %d", rc) clientMQTT = mqtt_client.Client(client_id) # clientMQTT.username_pw_set(username, password) clientMQTT.on_connect = on_connect clientMQTT.connect(broker, port) return clientMQTT def subscribe(clientMQTT: mqtt_client): def on_message(clientMQTT, userdata, msg): print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic") send_to_all_clients(msg.payload.decode()) clientMQTT.subscribe(topic) clientMQTT.on_message = on_message def runMQTT(): clientMQTT = connect_mqtt() subscribe(clientMQTT) clientMQTT.loop_start() # clientMQTT.loop_forever() def send_to_all_clients(message): # 获取当前Tornado的IO循环实例 io_loop = tornado.ioloop.IOLoop.current() # 将消息发送操作提交到IO循环线程执行 io_loop.add_callback(lambda: _send_to_all_clients(message)) def _send_to_all_clients(message): for clientws in clients: clientws.write_message(message) class SocketHandler(tornado.websocket.WebSocketHandler): def open(self): print("new client") clients.append(self) def on_close(self): clients.remove(self) print("removing client") def on_message(self, message): pass # for client in utility.clients: # if client != self: # client.write_message(msg) if __name__ == '__main__': app = tornado.web.Application( handlers = [ (r"/monitor", SocketHandler) ], debug = True, template_path = os.path.join(os.path.dirname(__file__), "templates"), static_path = os.path.join(os.path.dirname(__file__), "static") ) runMQTT() app.listen(8200) tornado.ioloop.IOLoop.instance().start()
关键修改点
- 拆分消息发送逻辑:新增
_send_to_all_clients负责实际的WebSocket消息发送 - 通过
IOLoop.current().add_callback()将消息发送操作切换到Tornado的IO循环线程中执行,确保WebSocket操作在正确的线程上下文运行
内容的提问来源于stack exchange,提问作者bb sof
相关产品推荐
相关产品推荐

