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
相关产品推荐
相关产品推荐

