如何在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
相关产品推荐
相关产品推荐

