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

如何在Tornado中终止旧循环并切换实时股票数据推送

解决Tornado WebSocket实时推送股票数据的循环阻塞问题

问题背景

基于Tornado开发的WebSocket服务,需实现以下功能:

  • 接收前端下拉菜单选中的股票代码
  • 从PostgreSQL获取对应股票的实时更新数据,推送给前端
  • 切换选中股票时,终止原股票的推送任务,转而推送新股票的数据

原代码在on_message方法中使用while True+time.sleep的阻塞循环,导致事件循环被卡住,无法接收新的前端消息,无法实现股票切换需求。

解决方案

核心思路是利用Tornado的异步事件循环机制,用周期性异步任务替代阻塞循环,同时维护任务标识以实现任务的终止与切换。

修改后的完整代码

import psycopg2
import os
import pandas as pd
import json
import tornado.websocket
import tornado.ioloop
from tornado import gen

class RealTimeDataHandler(tornado.websocket.WebSocketHandler):
    def open(self):
        print(f'Session Opened. IP:{self.request.remote_ip}')
        # 初始化数据库连接
        self.conn = psycopg2.connect(
            dbname='postgres',
            user='postgres',
            password=os.environ.get('PostgreSQL_PASS'),
            host='localhost',
        )
        self.cur = self.conn.cursor()
        # 追踪当前推送任务
        self.current_task = None
        self.ioloop = tornado.ioloop.IOLoop.current()

    def get_data_from_sql(self):
        self.cur.execute("SELECT * FROM public.random_time_multi")
        m = self.cur.fetchall()
        column_names = [desc[0] for desc in self.cur.description]
        data = pd.DataFrame(m, columns=column_names)
        return data

    @gen.coroutine
    def send_data_periodically(self, column_name):
        try:
            while True:
                # 获取并推送数据
                data = self.get_data_from_sql()
                message = json.dumps(data[column_name].to_dict())
                self.write_message(message)
                # 用Tornado异步超时替代time.sleep,避免阻塞事件循环
                yield gen.sleep(5)
        except tornado.websocket.WebSocketClosedError:
            # 连接关闭时捕获异常,避免报错
            pass

    def on_message(self, message):
        column_name = json.loads(message)['symbol']
        # 如果当前有推送任务,先取消
        if self.current_task is not None:
            self.current_task.cancel()
        # 启动新的推送任务
        self.current_task = self.ioloop.add_callback(self.send_data_periodically, column_name)

    def on_close(self):
        print("Session closed")
        # 关闭数据库资源
        self.cur.close()
        self.conn.close()
        # 取消当前推送任务
        if self.current_task is not None:
            self.current_task.cancel()
        
    def check_origin(self, origin):
        return True

def stop_tornado():
    tornado.ioloop.IOLoop.current().stop()

app = tornado.web.Application([(r"/ws/display", RealTimeDataHandler)])

if __name__ == '__main__':
    app.listen(8080)
    try:
        tornado.ioloop.IOLoop.current().start()
    except KeyboardInterrupt:
        print("KeyboardInterrupt received. Stopping Tornado.")
        stop_tornado()

关键修改说明

  • 新增任务追踪标识:self.current_task用来记录当前的推送任务,方便切换股票时取消旧任务
  • 异步周期性推送:将原阻塞循环改为send_data_periodically异步函数,用gen.sleep(5)替代time.sleep(5),不会阻塞Tornado的事件循环
  • 任务切换逻辑:在on_message中先取消当前任务(如果存在),再启动新的股票推送任务
  • 资源正确释放:on_close中关闭数据库连接并取消任务,原代码中self.ioloop.stop()会终止整个服务,改为只清理当前连接的资源
  • 异常处理:捕获WebSocketClosedError,避免连接关闭时抛出异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 23:02:48