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

如何实现常规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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 23:01:28