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

如何以非阻塞方式使用starmap_async?FastAPI场景解惑

解决方案:非阻塞使用starmap_async并避免FastAPI主线程阻塞

核心问题分析

你的代码阻塞主线程的原因有两个:

  1. starmap_async的wait()/get()是同步阻塞调用,在FastAPI的异步事件循环中执行会直接卡住整个线程,导致无法处理HTTP请求。
  2. 启动事件中直接await build_cache()会强制服务器等待缓存构建完成后才启动服务。

解决思路

将阻塞的多进程操作从异步事件循环中剥离,用asyncio.run_in_executor调度到独立线程/进程池执行;同时在启动事件中用后台任务异步执行缓存构建,不阻塞服务器启动。

代码实现方案一:基于multiprocessing.Pool改造

import asyncio
from multiprocessing.pool import Pool
from multiprocessing import cpu_count
from io import BytesIO
from ftplib import FTP
from fastapi import FastAPI
import uvicorn

ftp_url = "你的FTP地址"

def get_from_ftp(file_name):
    flo = BytesIO()
    ftp_server = FTP(ftp_url)
    ftp_server.login()
    ftp_server.retrbinary(f'RETR {file_name}', flo.write)
    flo.seek(0)
    # 写入磁盘逻辑
    with open(file_name, 'wb') as f:
        f.write(flo.read())

# 封装多进程下载逻辑为同步函数
def run_ftp_downloads():
    with Pool(cpu_count() - 1) as pool:
        some_files_to_download = ['file1', 'file2']
        # 单参数场景用map_async更合适,多参数才需要starmap_async(传入元组列表)
        result = pool.map_async(get_from_ftp, some_files_to_download)
        result.wait()

async def build_cache():
    loop = asyncio.get_running_loop()
    # 将阻塞的多进程操作放到executor中,不占用事件循环线程
    await loop.run_in_executor(None, run_ftp_downloads)

app = FastAPI()

@app.on_event("startup")
async def startup():
    # 用后台任务启动缓存构建,不等待完成,服务器直接启动
    asyncio.create_task(build_cache())

@app.get("/")
async def root():
    return {"message": "服务器已启动,缓存构建中..."}

if __name__ == "__main__":
    uvicorn.run(app, host="0.0.0.0", port=8000)

代码实现方案二:改用concurrent.futures.ProcessPoolExecutor(更贴合asyncio生态)

如果不需要multiprocessing.Pool的特殊API,推荐直接用concurrent.futures的进程池,和asyncio兼容性更好:

import asyncio
from concurrent.futures import ProcessPoolExecutor
from multiprocessing import cpu_count
from io import BytesIO
from ftplib import FTP
from fastapi import FastAPI
import uvicorn

ftp_url = "你的FTP地址"

def get_from_ftp(file_name):
    flo = BytesIO()
    ftp_server = FTP(ftp_url)
    ftp_server.login()
    ftp_server.retrbinary(f'RETR {file_name}', flo.write)
    flo.seek(0)
    with open(file_name, 'wb') as f:
        f.write(flo.read())

async def build_cache():
    loop = asyncio.get_running_loop()
    some_files_to_download = ['file1', 'file2']
    with ProcessPoolExecutor(max_workers=cpu_count()-1) as executor:
        # 将进程池的阻塞map操作放到executor中执行
        await loop.run_in_executor(None, lambda: list(executor.map(get_from_ftp, some_files_to_download)))

app = FastAPI()

@app.on_event("startup")
async def startup():
    asyncio.create_task(build_cache())

@app.get("/")
async def root():
    return {"message": "服务器已启动,缓存构建中..."}

if __name__ == "__main__":
    uvicorn.run(app, host="0.0.0.0", port=8000)

关键说明

  • 关于starmap_async的正确用法:如果get_from_ftp需要多个参数(比如get_from_ftp(file_name, ftp_user, ftp_pwd)),则需要将参数整理为元组列表,用starmap_async,例如:
    params = [('file1', 'user', 'pwd'), ('file2', 'user', 'pwd')]
    result = pool.starmap_async(get_from_ftp, params)
    
  • Windows系统注意事项:多进程代码必须放在if __name__ == "__main__"保护之外,否则会触发进程启动异常。

内容的提问来源于stack exchange,提问作者jossefaz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 10:05:23