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

如何在子线程中监听事件?实现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接口调用,直接通过消息频道传递控制指令。

具体实现思路

  1. 编写异步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
    
  2. 整合到现有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()
    
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 11:33:25