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

如何使用Python async/await同时处理多个WebSocket消息?

解决WebSocket消息监听与耗时任务并发执行的问题

嘿,我来帮你搞定这个问题!你踩的这个坑其实是async/await里很常见的——当你用await直接调用on_message()时,当前的协程会暂停下来,直到这个耗时数分钟的任务完全结束,自然就没法继续接收新的WebSocket消息了。

核心思路很简单:把on_message()包装成独立的后台任务,让它和WebSocket监听协程同时运行,不用等它完成就能继续接收下一条消息。下面是修改后的完整代码,我会逐段解释变化:

import websockets
import asyncio

my_socket = "ws://......."

async def on_message(message):
    # 给耗时任务加上异常捕获,避免后台任务出错时无提示
    try:
        # Do what needs to be done with received message
        # 这里保留你的耗时操作,比如带await的sleep都没问题
        await asyncio.sleep(120)  # 模拟2分钟的处理流程
        print(f"✅ 处理完成消息: {message[:20]}...")
    except Exception as e:
        print(f"❌ 处理消息出错: {str(e)}")

async def listen_websocket():
    while True:
        try:
            async with websockets.connect(my_socket) as ws:
                print("🔌 WebSocket连接成功,开始监听消息")
                while True:
                    message = await ws.recv()
                    # 关键改动:用create_task启动后台任务,无需等待完成
                    asyncio.create_task(on_message(message))
                    print(f"📥 收到新消息,已启动处理任务: {message[:20]}...")
        except Exception as e:
            print(f"⚠️ WebSocket连接异常: {str(e)}")
            # 等待10秒后重连,避免频繁请求服务器
            await asyncio.sleep(10)

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

关键改动说明

  1. 拆分接收与处理逻辑:
    原来的await on_message(await ws.recv())会让监听流程完全阻塞,现在拆成两步:先接收消息,再用asyncio.create_task()把on_message()变成一个独立的后台任务。这样监听协程可以立刻回到await ws.recv()继续等待新消息,两者互不干扰。

  2. 给on_message()加异常捕获:
    后台任务的异常默认不会主动抛出(只会存在于任务的result()中),加上try-except能让你及时发现处理消息时的错误,避免“无声失败”。

  3. 结构化代码:
    把外层的重连逻辑封装成listen_websocket()协程,用asyncio.run()启动整个事件循环,代码结构更清晰易维护。

额外注意事项

  • 如果你的on_message()里有同步阻塞代码(比如time.sleep(1)而不是await asyncio.sleep(1)),哪怕用create_task也会阻塞整个事件循环。这时候可以用asyncio.to_thread()把同步代码放到线程池里运行:
    import time
    async def on_message(message):
        # 把同步阻塞代码放到线程中执行
        await asyncio.to_thread(lambda: time.sleep(120))
    
  • 如果需要限制并发处理的任务数量(比如怕太多任务占用资源),可以用asyncio.Semaphore做流量控制:
    # 定义信号量,最多同时处理5个任务
    semaphore = asyncio.Semaphore(5)
    
    async def on_message(message):
        async with semaphore:
            # 这里是你的耗时操作
            await asyncio.sleep(120)
            print(f"✅ 处理完成消息: {message[:20]}...")
    

这样就能完美实现你的需求:持续监听WebSocket消息,同时并行处理每条消息的耗时任务,再也不会出现后续消息延迟的问题啦!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 09:34:07