Python中实现类似Promise.all的异步并发调用问题
问题描述
我想在Python中实现类似JavaScript await Promise.all() 的并发功能,找到了asyncio.gather(),但实际用下来没实现异步执行。我的场景是基于FastAPI从S3的多个远程文件中提取指定值,全部完成后收集结果——同样的任务在JS里处理12个文件只需要1秒多,但现在Python代码里读取的文件越多耗时越长,明显是串行执行的。我怀疑是rasterio流式读取远程文件的操作没法异步,想知道怎么改代码才能让函数并发调用,最后统一收集响应。
简化后的代码如下:
async def read_from_file(s3_path): # 注意:这里是通过S3路径流式读取远程文件 with rasterio.open(s3_path) as src: values = src.read(1, window=Window(1, 2, 1, 1)) return values[0][0] @app.get("/get-all") async def get_all(): start_time = datetime.datetime.now() # 示例路径 s3_paths = [ "s3:file-1", "s3:file-2", "s3:file-3", "s3:file-4", "s3:file-5", "s3:file-6", ] values = await asyncio.gather( read_from_file(s3_paths[0]), read_from_file(s3_paths[1]), read_from_file(s3_paths[2]), read_from_file(s3_paths[3]), read_from_file(s3_paths[4]), read_from_file(s3_paths[5]), ) end_time = datetime.datetime.now() logger.info(f"duration: {end_time-start_time}")
问题根源
你代码里的read_from_file虽然定义成了async函数,但内部的rasterio.open()和src.read()都是同步阻塞的IO操作。在Python的asyncio事件循环中,只要有一个异步函数里存在同步阻塞代码,就会卡住整个事件循环,导致所有任务只能串行执行,这就是为什么文件越多耗时越长。
解决方案
把同步阻塞的IO操作放到线程池中执行,让asyncio事件循环可以同时调度多个任务,实现真正的并发。推荐用asyncio.to_thread()(Python 3.9+支持),它能方便地把同步函数包装成可异步调用的任务。
修改后的代码
import asyncio import datetime import rasterio from rasterio.windows import Window from fastapi import FastAPI import logging logger = logging.getLogger(__name__) app = FastAPI() # 把同步读取逻辑封装成普通函数 def sync_read_from_file(s3_path): with rasterio.open(s3_path) as src: values = src.read(1, window=Window(1, 2, 1, 1)) return values[0][0] async def read_from_file(s3_path): # 用to_thread把同步操作放到线程池执行 return await asyncio.to_thread(sync_read_from_file, s3_path) @app.get("/get-all") async def get_all(): start_time = datetime.datetime.now() s3_paths = [ "s3:file-1", "s3:file-2", "s3:file-3", "s3:file-4", "s3:file-5", "s3:file-6", ] # 生成所有异步任务,用gather并发执行 tasks = [read_from_file(path) for path in s3_paths] values = await asyncio.gather(*tasks) end_time = datetime.datetime.now() logger.info(f"duration: {end_time-start_time}") return {"values": values}
关键说明
asyncio.to_thread()会把同步函数提交到默认的线程池,每个任务在独立线程中执行,不会阻塞asyncio的事件循环,这样多个文件读取任务就能同时进行。- 如果你的Python版本低于3.9,可以用
concurrent.futures.ThreadPoolExecutor手动创建线程池,用法类似:from concurrent.futures import ThreadPoolExecutor executor = ThreadPoolExecutor(max_workers=10) # 可根据文件数量调整线程数 async def read_from_file(s3_path): loop = asyncio.get_running_loop() return await loop.run_in_executor(executor, sync_read_from_file, s3_path) - 线程数不需要设置得过大,S3的并发请求有上限,一般设置为10-20就能达到最优效果。
内容的提问来源于stack exchange,提问作者sobmortin354
相关产品推荐
相关产品推荐

