如何在子线程中监听事件?实现Binance WebSocket跨应用控制
问题:Redis Pub/Sub是否适合实现Django对独立FastAPI WebSocket服务的控制?
我正在开发基于Binance WebSocket API的信息获取应用,已有Django主代码库,将WebSocket模块实现为独立的FastAPI应用,代码如下:
import asyncio from threading import Thread import uvicorn from fastapi import FastAPI from binance import AsyncClient, BinanceSocketManager app = FastAPI() @app.get("/") async def root(): return "Websocket app" async def main_socket_stream(): client = await AsyncClient.create() binance_manager = BinanceSocketManager(client=client) multiplex_socket = binance_manager.futures_multiplex_socket(symbols) async with multiplex_socket as active_socket: while True: result = await active_socket.recv() print(result) # There are lots of further logic - proceed data, save something to db, etc. def side_thread(): loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) asyncio.run(main_socket_stream()) if __name__ == "__main__": thread = Thread(target=side_thread, args=(), daemon=True) thread.start() uvicorn.run(app, port=5105)
目前后台WebSocket已正常运行,但我希望能从Django主应用对其进行控制,比如停止WebSocket、重载更新后的交易对列表等。我设想通过如下逻辑实现:
async with multiplex_socket as active_socket: while True: if await EventListenerClass.new_event_fired(): break # or do something special result = await active_socket.recv() print(result)
考虑使用Redis Pub/Sub实现一个EventListener类,监听指定频道的消息(如"stop_websocket"、"reload_with_updated_symbols_list"等)。我需要满足以下需求的解决方案:
- 响应速度足够快(每秒级)
- 支持对WebSocket的控制
请问使用Redis的方案是否可行?
回答
Redis Pub/Sub方案完全可行,完全能满足你的需求:
为什么可行?
- 响应速度达标:Redis Pub/Sub的消息推送是低延迟的,完全能达到你要求的每秒级响应,甚至更快,不会成为控制逻辑的瓶颈。
- 跨服务通信适配:你的Django主应用和FastAPI WebSocket服务是独立部署的,Redis Pub/Sub刚好能作为轻量级的跨进程/跨服务通信中间件,不需要复杂的HTTP接口调用,直接通过消息频道传递控制指令。
具体实现思路
编写异步Redis Pub/Sub监听类
用异步Redis客户端(比如aioredis)实现EventListener类,在单独的异步任务里监听指定频道,收到消息后触发对应的事件标记:import aioredis import asyncio class EventListener: def __init__(self, redis_url, channel): self.redis_url = redis_url self.channel = channel self.redis = None self.event_queue = asyncio.Queue() self.stop_flag = False self.reload_flag = False self.symbols = [] async def connect(self, initial_symbols): self.symbols = initial_symbols self.redis = await aioredis.from_url(self.redis_url) asyncio.create_task(self._listen()) async def _listen(self): async with self.redis.pubsub() as pubsub: await pubsub.subscribe(self.channel) async for message in pubsub.listen(): if message["type"] == "message": await self.event_queue.put(message["data"].decode()) async def new_event_fired(self): # 非阻塞检查队列是否有消息,避免阻塞WebSocket接收逻辑 try: event = self.event_queue.get_nowait() self._handle_event(event) return True except asyncio.QueueEmpty: return False def _handle_event(self, event): # 处理不同的控制指令,比如设置停止标记、更新交易对列表 if event == "stop_websocket": self.stop_flag = True elif event.startswith("reload_symbols:"): self.symbols = event.split(":")[1].split(",") self.reload_flag = True整合到现有WebSocket逻辑
修改main_socket_stream函数,初始化EventListener并在循环里检查事件:# 假设初始交易对列表定义在这里 INITIAL_SYMBOLS = ["BTCUSDT", "ETHUSDT"] async def main_socket_stream(): # 初始化事件监听器 event_listener = EventListener("redis://localhost:6379", "websocket_control") await event_listener.connect(INITIAL_SYMBOLS) client = await AsyncClient.create() binance_manager = BinanceSocketManager(client=client) multiplex_socket = binance_manager.futures_multiplex_socket(event_listener.symbols) async with multiplex_socket as active_socket: while not event_listener.stop_flag: # 先检查控制事件 await event_listener.new_event_fired() # 如果需要重载交易对 if event_listener.reload_flag: # 关闭当前socket,重新创建新的multiplex socket await active_socket.close() multiplex_socket = binance_manager.futures_multiplex_socket(event_listener.symbols) active_socket = await multiplex_socket.__aenter__() event_listener.reload_flag = False # 接收Binance数据 result = await active_socket.recv() print(result) # 其他业务逻辑... # 清理资源 await client.close()Django端发送控制指令
在Django里用同步Redis客户端(比如redis-py)向指定频道发送消息:import redis def send_websocket_control(event): r = redis.Redis(host='localhost', port=6379, db=0) r.publish("websocket_control", event) # 示例:停止WebSocket send_websocket_control("stop_websocket") # 示例:重载交易对 send_websocket_control("reload_symbols:BTCUSDT,ETHUSDT,BNBUSDT")
注意事项
- 确保Redis服务稳定运行,这是跨服务通信的核心。
- 异步Redis客户端要和FastAPI的事件循环兼容,避免阻塞。
- 控制指令的格式要统一,方便EventListener解析处理。
内容的提问来源于stack exchange,提问作者Vladislav
相关产品推荐
相关产品推荐

