You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 08:15:43