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

WebSocket多流场景下,代码执行时长超推送间隔时所有消息都会被处理吗?

现有代码的基础问题

你当前写的代码首先无法实现多路流接收的需求:websockets.connect() 单次仅支持传入单个WebSocket服务端URL,直接传入URL列表会直接抛出参数类型错误,你需要为每个流地址单独创建对应的连接协程。

同步处理逻辑阻塞的运行机制

假设你已经修正了多路连接的实现问题,为每个流单独运行处理协程,那么处理逻辑的耗时影响分两种情况:

  • 如果你的处理逻辑是同步阻塞代码(比如CPU密集计算、同步HTTP请求、同步文件读写等),会直接卡住整个asyncio事件循环:当14:05收到的stream1数据处理耗时70秒时,这70秒内整个事件循环完全被占住,所有其他WebSocket连接的接收逻辑、其他协程的执行都会被暂停。
  • 如果你的处理逻辑是全异步实现(所有IO操作都通过async/await调用异步库实现),那么处理过程中事件循环会自动切换到其他协程处理stream2的数据,不会出现阻塞问题,仅处理逻辑中的CPU密集运算部分会短暂卡住事件循环。

同批次数据的处理规则

只要服务端推送的14:05批次stream2的数据没有因为内核接收缓冲区满被丢弃,所有数据都会按接收顺序排队,不会直接跳过14:05的内容去处理14:06的批次,只会整体延后执行,所有积压的数据都会依次处理,不会主动丢弃。

线程优化的必要性

你这种仅小概率出现处理超时、不会产生大量积压的场景,不需要做复杂的线程池改造,仅需用Python 3.9+内置的asyncio.to_thread()将同步处理逻辑丢到独立线程运行即可,改造后的参考代码如下:

import asyncio
import json
import websockets

stream_urls = [stream_url1, stream_url2, etc.]
reply_timeout = 60

# 单个流的处理协程
async def handle_stream(url):
    async with websockets.connect(url, ping_timeout=None) as sock:
        while True:
            res = await asyncio.wait_for(sock.recv(), timeout=reply_timeout)
            res = json.loads(res)
            # 将同步处理逻辑丢到线程执行,不阻塞事件循环
            await asyncio.to_thread(你的处理函数, res)

# 批量启动所有流的处理协程
async def main():
    stream_tasks = [handle_stream(url) for url in stream_urls]
    await asyncio.gather(*stream_tasks)

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

这种改造成本极低,也不会带来额外的性能开销,可以完全避免小概率的处理阻塞影响其他流的正常接收执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 11:06:05