使用asyncio实现双层异步时subprocess调用未并行的问题
问题:两层异步并行逻辑未生效的原因
输入数据结构如下:
input = { "set1": ["subset11", "subset12"], "set2": ["subset21", "subset22"], ... }
需求:
- 并行处理所有set
- 每个set内并行处理subset,之后进行汇总
编写的代码如下:
# top level function async def process_all_sets(input): tasks = [process_set(payload) for key, payload in input.items()] results = await asyncio.gather(*tasks) # process a single set async def process_set(payload): tasks = [process_subset(item) for item in payload] results = await asyncio.gather(*tasks) # here, loop over results and do some summarization # and return it return summary # process a single subset async def process_subset(subset): # need to run a subprocess here, it may take several minutes subprocess.run("some_command_based_on_subset") # do whatever needs to be done after subprocess completes # and return result return result
按照asyncio的设计,同一set内的多个process_subset调用应并行执行,期望多个subprocess.run能同时启动,但实际每次仅能看到一个调用执行,请问并行性未生效的原因是什么?
原因与解决方案
核心原因
你代码里的subprocess.run()是同步阻塞调用,它会直接卡住当前的asyncio事件循环。因为asyncio基于单线程事件循环运行,当某个协程调用同步阻塞函数时,事件循环会被完全占用,无法切换到其他协程执行,导致所有process_subset只能串行执行,无法实现并行。
解决方案
替换同步的subprocess.run为asyncio提供的异步子进程API,比如asyncio.create_subprocess_exec或asyncio.create_subprocess_shell,这些API不会阻塞事件循环,能让多个子进程真正并行启动。
修改后的process_subset示例:
async def process_subset(subset): # 使用异步子进程API proc = await asyncio.create_subprocess_shell( f"some_command_based_on_subset {subset}", stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE ) # 等待子进程完成并获取输出 stdout, stderr = await proc.communicate() # 后续处理逻辑 return {"stdout": stdout.decode(), "stderr": stderr.decode()}
这样修改后,当asyncio.gather调度多个process_subset协程时,每个协程在调用await proc.communicate()时会主动让出事件循环,让其他协程有机会启动各自的子进程,从而实现真正的并行。
内容的提问来源于stack exchange,提问作者shikhanshu
相关产品推荐
相关产品推荐

