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

Tornado WebSocket结合设备数据读取循环的最佳实践是什么?

Tornado WebSocket 持续推送设备数据的最佳实践

嘿,你遇到的这个问题其实是Tornado新手常踩的典型坑——把阻塞的无限循环直接塞进了open方法里!Tornado是单线程异步框架,open被这个循环卡住后,整个事件循环(IOLoop)就没法处理任何其他事件了,比如新的连接请求、客户端关闭动作,甚至其他异步任务都得排队等着,这就是为啥你的on_close和新连接都没反应。

下面是业界通用的解决方案,我一步步给你拆解:

1. 维护活跃连接集合

我们需要一个地方来跟踪所有当前在线的WebSocket客户端,这样才能方便地给所有人推送数据。这里用类变量来存储最方便,所有实例都能访问到。

2. 用定时任务替代阻塞循环

别再用while True死循环了,Tornado提供了PeriodicCallback(或者Tornado 5.0+推荐的add_callback_periodic)来定期执行数据读取和推送操作,这样不会阻塞事件循环,其他事件该怎么处理就怎么处理。

如果你的设备读取操作是阻塞且耗时的(比如要等设备响应几百毫秒),那还得结合线程池来避免卡住IOLoop,我后面会补充这种场景的处理方式。

修改后的完整代码

#!/usr/bin/env python
import tornado.httpserver
import tornado.websocket
import tornado.ioloop
import tornado.web
import socket

class MyWebSocketServer(tornado.websocket.WebSocketHandler):
    # 类变量:存储所有活跃的WebSocket连接
    active_clients = set()

    def open(self):
        client_ip = self.request.remote_ip
        print(f'New connection from {client_ip}')
        # 新连接加入集合
        self.active_clients.add(self)

    def on_close(self):
        client_ip = self.request.remote_ip
        print(f'Connection closed from {client_ip}')
        # 连接关闭时从集合移除
        self.active_clients.remove(self)

    @classmethod
    def broadcast_data(cls, data):
        # 遍历所有客户端推送数据,注意要处理连接已关闭的情况
        # 用list()转一下,避免遍历中集合被修改引发异常
        for client in list(cls.active_clients):
            try:
                client.write_message(data)
            except tornado.websocket.WebSocketClosedError:
                # 遇到已关闭的连接,直接从集合里清掉
                cls.active_clients.remove(client)
                print(f"Cleaned up a closed connection")

def read_device_data():
    # 这里替换成你实际读取设备数据的逻辑
    try:
        # 示例:模拟生成设备数据
        data = f"Device update: {socket.gethostname()} | Timestamp: {tornado.ioloop.IOLoop.current().time()}"
        return data
    except Exception as error:
        print(f"Failed to read device data: {str(error)}")
        return None

def push_data_to_clients():
    # 读取设备数据
    device_data = read_device_data()
    if device_data:
        # 广播给所有在线客户端
        MyWebSocketServer.broadcast_data(device_data)

application = tornado.web.Application([(r'/ws', MyWebSocketServer),])

if __name__ == "__main__":
    http_server = tornado.httpserver.HTTPServer(application)
    http_server.listen(8000)
    print('Server running on port 8000')
    
    # 初始化定时任务:每1秒推送一次数据(可根据你的需求调整间隔,单位毫秒)
    data_push_callback = tornado.ioloop.PeriodicCallback(push_data_to_clients, 1000)
    data_push_callback.start()
    
    tornado.ioloop.IOLoop.instance().start()

关键细节说明

  • active_clients集合:用类变量存储所有活跃连接,open时添加、on_close时移除。遍历的时候转成list()是为了防止遍历过程中集合被修改(比如移除已关闭的连接)导致迭代异常。
  • broadcast_data类方法:封装了广播逻辑,同时处理了WebSocketClosedError异常,确保集合里始终是有效的活跃连接。
  • PeriodicCallback:定时触发数据读取和推送,间隔时间可以根据设备数据的更新频率灵活调整。
  • 处理阻塞的设备读取:如果你的read_device_data是阻塞且耗时的,别直接在定时任务里调用,得把它放到线程池里异步执行,示例代码如下:
    async def async_read_and_push():
        loop = tornado.ioloop.IOLoop.current()
        # 把阻塞的读取操作交给线程池处理,不卡主事件循环
        device_data = await loop.run_in_executor(None, read_device_data)
        if device_data:
            MyWebSocketServer.broadcast_data(device_data)
    
    # 启动定时任务的方式也要调整(Tornado 5.0+推荐)
    async def main():
        http_server.listen(8000)
        print('Server running on port 8000')
        # 每隔1秒执行一次异步推送任务
        tornado.ioloop.IOLoop.current().add_callback_periodic(async_read_and_push, 1)
        await tornado.ioloop.IOLoop.current().future()
    
    if __name__ == "__main__":
        tornado.ioloop.IOLoop.current().run_sync(main)
    

这样修改后,你的服务器既能持续给所有客户端推送设备数据,又能正常处理新连接和关闭请求,完全符合Tornado异步框架的最佳实践。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:16:46