使用aiohttp后台下载大文件时CPU密集任务致下载暂停的问题排查
问题描述
我需要下载一系列约200MB的大文件,希望利用下载时间进行CPU密集型处理。我正在研究asyncio和aiohttp,原本认为可以用它们启动大文件下载后,在同一线程中后台继续下载的同时执行繁重计算。
但实际发现,在CPU密集任务运行期间下载会暂停,计算完成后才恢复。我监控了脚本运行时的CPU和带宽,明确看到约30秒的计算过程中下载暂停。以下是最小复现示例:
import asyncio import time import aiofiles import aiohttp async def download(session): url = 'https://repo.anaconda.com/archive/Anaconda3-2022.10-Linux-s390x.sh' # 280 MB file async with session.get(url) as resp: async with aiofiles.open('./tmpfile', mode='wb') as f: print('Starting the download') data = await resp.read() print('Starting the file write') await f.write(data) print('Download completed') async def heavy_cpu_load(): await asyncio.sleep(5) # make sure the download has started print('Starting the computation') for i in range(200000000): # takes about 30 seconds on my laptop. i ** 0.5 print('Finished the computation') async def main(): async with aiohttp.ClientSession() as session: timer = time.time() tasks = [download(session), heavy_cpu_load()] await asyncio.gather(*tasks) print(f'All tasks completed in {time.time() - timer}s') if __name__ == '__main__': asyncio.run(main())
请问我是操作有误,还是对aiohttp的能力理解有误?
问题原因与解决方案
核心原因:asyncio是单线程事件循环模型,它的任务调度完全依赖
await关键字——只有当任务执行到await(等待IO操作、sleep等可挂起的操作)时,事件循环才会切换到其他任务。而你的heavy_cpu_load函数里的for循环没有任何await,一旦开始执行就会霸占整个线程,事件循环根本没机会调度下载任务,导致下载暂停。操作修正方案:把CPU密集型任务放到单独的线程或进程中执行,避免阻塞事件循环。推荐使用Python 3.9+提供的
asyncio.to_thread(),它可以轻松把同步函数提交到线程池运行:
修改后的代码示例:
import asyncio import time import aiofiles import aiohttp # 把CPU密集计算抽成单独的同步函数 def compute(): for i in range(200000000): i ** 0.5 async def download(session): url = 'https://repo.anaconda.com/archive/Anaconda3-2022.10-Linux-s390x.sh' # 280 MB file async with session.get(url) as resp: async with aiofiles.open('./tmpfile', mode='wb') as f: print('Starting the download') # 流式读取大文件,增加await调度点同时降低内存占用 async for chunk in resp.content.iter_chunked(1024*1024): await f.write(chunk) print('Download completed') async def heavy_cpu_load(): await asyncio.sleep(5) # 确保下载启动 print('Starting the computation') # 把计算任务提交到线程池,不阻塞事件循环 await asyncio.to_thread(compute) print('Finished the computation') async def main(): async with aiohttp.ClientSession() as session: timer = time.time() tasks = [download(session), heavy_cpu_load()] await asyncio.gather(*tasks) print(f'All tasks completed in {time.time() - timer}s') if __name__ == '__main__': asyncio.run(main())
- 额外优化:原下载代码用
resp.read()一次性读取280MB文件到内存,改成流式分块读取(resp.content.iter_chunked)不仅能降低内存占用,还能增加更多await调度点,让事件循环在下载过程中更灵活地切换任务。
内容的提问来源于stack exchange,提问作者rowan_uk
相关产品推荐
相关产品推荐

