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

Python使用websockets与asyncio监听多Socket失效问题及解决方案咨询

解决多Websocket监听问题的方案

嘿,我来帮你搞定这个问题!你遇到的“只能监听第一个Socket”的情况,大概率是因为你在处理第一个Socket的无限循环时,没有给其他异步任务执行的机会。Asyncio是协作式多任务,只要你的代码里有await(比如websocket的接收操作),就不会完全阻塞,但如果你的任务启动方式不对,还是会出现“卡住第一个”的情况。下面给你两种可行的解决办法:

优先方案:用Asyncio原生方式管理多任务

这是最符合Asyncio设计理念的方案,不需要额外引入线程,性能也更好。核心思路是:把所有Websocket的订阅任务都创建好,然后用asyncio.gather()让它们同时运行。

修正后的代码示例

import asyncio
import json
import websockets

class SocketManager:
    def __init__(self):
        self.tasks = []

    async def subscribe(self, event):
        # 替换成你的实际Websocket地址
        ws_uri = "ws://your-target-websocket-server"
        async with websockets.connect(ws_uri) as websocket:
            # 发送订阅请求
            payload = json.dumps(event)
            await websocket.send(payload)
            # 无限循环监听消息——这里的await是关键,会让出CPU给其他任务
            while True:
                message = await websocket.recv()
                print(f"[{event['channel']}] 收到消息: {message}")

    def add_socket_task(self, event):
        # 只创建任务,不立即执行,避免阻塞后续任务添加
        current_loop = asyncio.get_event_loop()
        self.tasks.append(current_loop.create_task(self.subscribe(event)))

    async def run_all_sockets(self):
        # 同时运行所有已添加的Socket任务
        await asyncio.gather(*self.tasks)

# 用法示例
if __name__ == "__main__":
    manager = SocketManager()
    # 添加多个Socket监听任务
    manager.add_socket_task({"channel": "trade_btc"})
    manager.add_socket_task({"channel": "trade_eth"})
    manager.add_socket_task({"channel": "order_book"})
    # 启动所有任务
    asyncio.run(manager.run_all_sockets())

为什么这个方案能解决问题?

  • 你之前的代码可能是在调用subscribe时直接await了,导致第一个无限循环占用了整个事件循环,后面的任务根本没机会启动。
  • 现在我们先把所有任务都创建好存入列表,最后用asyncio.gather()一次性启动所有任务。每个subscribe里的await websocket.recv()会主动让出CPU,让其他任务也能执行。

备选方案:为每个Socket单独使用线程

如果你的场景中存在无法异步化的同步阻塞代码(比如某些老库不支持Asyncio),可以给每个Socket单独开一个线程,每个线程运行自己的Asyncio事件循环。

线程方案的代码示例

import asyncio
import json
import websockets
import threading

def run_single_socket(event):
    async def subscribe():
        ws_uri = "ws://your-target-websocket-server"
        async with websockets.connect(ws_uri) as websocket:
            payload = json.dumps(event)
            await websocket.send(payload)
            while True:
                message = await websocket.recv()
                print(f"[{event['channel']}] 收到消息: {message}")
    
    # 每个线程创建独立的事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    loop.run_until_complete(subscribe())

# 用法示例
if __name__ == "__main__":
    # 为每个Socket创建线程
    thread1 = threading.Thread(target=run_single_socket, args=({"channel": "trade_btc"},))
    thread2 = threading.Thread(target=run_single_socket, args=({"channel": "trade_eth"},))
    
    thread1.start()
    thread2.start()
    
    # 等待所有线程结束(可选,根据你的需求)
    thread1.join()
    thread2.join()

注意事项

  • 线程间的资源共享需要注意加锁(比如如果多个线程要操作同一个全局变量),但如果只是各自监听独立的Socket,基本不需要额外处理。
  • 这种方案的资源开销比纯Asyncio大,所以如果没有特殊需求,优先用第一种方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:24:58