Python asyncio异步任务未按预期并行执行问题排查
问题原因
代码没有按预期并发,核心有三个问题:
- CPU密集逻辑阻塞asyncio事件循环:asyncio是单线程协作式并发模型,只有任务执行到
await关键字主动让出控制权时,事件循环才能调度其他任务运行。你写的10亿次累加循环是纯CPU密集运算,全程没有任何让出控制权的操作,单个任务启动后会一直霸占CPU线程直到执行完毕,其余任务完全没有运行机会,自然无法并发。 - 任务调度顺序不符合预期:你在循环中调用
asyncio.create_task只是把任务加入就绪队列,等所有任务创建完成、代码执行到await asyncio.gather时才会开始调度任务。第一个任务被调度后,先打印chunk_no:0,紧接着就进入阻塞式的CPU计算,直到整个分块处理完才会把控制权交还给事件循环,第二个任务这时候才能启动打印自己的chunk_no,就出现了你看到的串行输出结果。 - 分块索引计算错误:计算分块起始位置时你用了
i * self.thread_count,正确逻辑应该是i * self.chunk_size,当前写法会导致分块范围错乱,出现数据重复读取、遗漏的问题。
修复方法
要实现「先打印所有分块编号,再并行处理数据」的效果,按以下逻辑调整:
- 拆分打印和计算逻辑:打印chunk_no后主动加一次
await asyncio.sleep(0)强制让出控制权,保证所有任务都能先完成编号打印。 - CPU密集运算不要直接跑在事件循环线程中:用进程池承载CPU计算逻辑,既可以避免阻塞事件循环,还能绕过Python GIL限制,利用多核实现真正的并行加速。
- 修正分块索引的计算逻辑。
修复后的核心代码如下:
import asyncio from concurrent.futures import ProcessPoolExecutor # 抽离CPU密集计算逻辑为普通同步函数,交给进程池执行 def _process_chunk_sync(chunk_no, chunk): for i, seq in enumerate(chunk): counter = 0 for j in range(i, 1000000000): counter += 1 print('Processed %d. item in chunk %d' % (i, chunk_no)) class Extractor: async def process_chunk(self, chunk_no, chunk, executor): print(f'chunk_no: {chunk_no}') # 主动让出控制权,保证其他任务的打印逻辑能优先执行 await asyncio.sleep(0) loop = asyncio.get_running_loop() # 将CPU计算丢到进程池执行,不阻塞事件循环 await loop.run_in_executor(executor, _process_chunk_sync, chunk_no, chunk) async def run(self): ctx = Wtp() ctx.process("~/Dev/wikipedia/enwiki-latest-pages-articles15.xml-p17324603p17460152.bz2", page_handler) self.thread_count = 256 self.chunk_size = int(len(ctx.page_seq) / self.thread_count) tasks = [] # 初始化进程池,进程数和分块数匹配即可 executor = ProcessPoolExecutor(max_workers=self.thread_count) for i in range(0, self.thread_count): # 修复索引计算错误 fromIndex = i * self.chunk_size toIndex = fromIndex + self.chunk_size if fromIndex + self.chunk_size < len(ctx.page_seq) else len(ctx.page_seq) chunk = ctx.page_seq[fromIndex:toIndex] tasks.append(asyncio.create_task( self.process_chunk(i, chunk, executor))) await asyncio.gather(*tasks) executor.shutdown() async def main(): extractor = Extractor() await extractor.run() if __name__ == '__main__': asyncio.run(main())
补充说明:
await asyncio.sleep(0)不会产生实际的等待,只会触发一次事件循环调度,让所有排队的任务都有机会执行到打印chunk_no的逻辑,就能实现先输出所有分块编号、再开始处理的效果。- 如果你的处理逻辑是IO密集型(比如网络请求、文件读写),不需要用进程池,只要在IO操作的位置调用对应异步方法加
await即可;但纯CPU计算必须放到独立的执行单元中,否则asyncio永远无法实现并发。
内容的提问来源于stack exchange,提问作者cuneyttyler
相关产品推荐
相关产品推荐

