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

Asyncio同时维护两个WebSocket连接时无法执行第三个异步函数

问题分析与解决方案

核心错误点

  1. 事件循环被阻塞:如果你的WebSocket连接用了同步式死循环(比如while True持续收消息但不让出事件循环),会直接霸占异步事件循环,导致checking()完全没机会执行。异步编程的核心是任务间要主动让出执行权,否则单个任务会占用所有资源。
  2. checking()未被正确注册到事件循环:你可能只定义了checking()函数,但没通过asyncio.create_task()这类方法把它加入异步事件循环;或者启动顺序错误,导致它根本没启动。
  3. 全局变量的异步适配问题:直接用全局变量传递数据时,checking()如果没有合理的触发逻辑(比如异步通知而非定时轮询),要么拿不到新数据,要么轮询逻辑本身又阻塞了事件循环。

正确实现思路

放弃全局变量,改用异步队列(asyncio.Queue)传递数据,同时确保所有异步任务都被正确注册到事件循环并行运行:

  1. 两个WebSocket任务负责接收数据,将数据放入队列;
  2. checking()任务持续从队列取数据处理;
  3. 用asyncio.create_task()启动所有任务,通过asyncio.gather()维持任务运行。

示例代码

import asyncio
import websockets

async def kraken_websocket(queue):
    # 连接Kraken WebSocket
    uri = "wss://ws.kraken.com"
    async with websockets.connect(uri) as ws:
        # 订阅BTC/USDT行情(示例订阅指令)
        await ws.send('{"event":"subscribe","pair":["XBT/USD"],"subscription":{"name":"ticker"}}')
        # 异步接收消息,自动让出事件循环
        async for msg in ws:
            await queue.put(("kraken", msg))

async def binance_websocket(queue):
    # 连接Binance WebSocket
    uri = "wss://stream.binance.com:9443/ws/btcusdt@ticker"
    async with websockets.connect(uri) as ws:
        async for msg in ws:
            await queue.put(("binance", msg))

async def checking(queue):
    while True:
        # 异步等待队列新数据,不会阻塞事件循环
        source, data = await queue.get()
        # 这里写你的数据处理逻辑
        print(f"从{source}收到数据:{data[:60]}...")
        # 标记任务完成(可选,用于队列join()方法)
        queue.task_done()

async def main():
    # 创建异步队列用于数据传递
    data_queue = asyncio.Queue()
    # 启动三个异步任务
    kraken_task = asyncio.create_task(kraken_websocket(data_queue))
    binance_task = asyncio.create_task(binance_websocket(data_queue))
    checking_task = asyncio.create_task(checking(data_queue))
    # 等待所有任务持续运行(手动终止前不会退出)
    await asyncio.gather(kraken_task, binance_task, checking_task)

if __name__ == "__main__":
    asyncio.run(main())

关键说明

  • async for处理WebSocket消息:这种方式会在每次接收消息后主动让出事件循环,让checking()有执行机会;
  • 异步队列替代全局变量:队列是异步安全的,能避免全局变量的竞态问题,数据传递更可靠;
  • 任务并行启动:通过create_task()将三个任务都加入事件循环,由asyncio自动调度并行运行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 22:07:12