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

如何优化FastAPI端点的Python并发请求?提升外部分页API获取效率

问题分析与优化方案

当前实现的潜在问题

  • 线程安全隐患:parameters是共享字典,多线程同时修改parameters["page"]会导致请求参数混乱,比如多个请求可能携带错误页码。必须每次创建参数副本,而非修改同一字典。
  • 同步IO低效:requests是同步HTTP库,线程池靠多线程掩盖IO等待,但线程有创建、切换开销,IO密集型场景下异步IO效率更高。
  • 数据处理错误:items += data直接拼接requests.Response对象,而非解析后的业务数据,会导致列表存的是响应对象而非实际内容,应改为解析响应(如response.json())后再添加。
  • 并发数不合理:CONNECTIONS=100远大于实际请求数(15-25),过多线程只会增加系统开销,甚至触发外部API限流,拖慢整体速度。

修复后的线程池实现

先解决核心问题,优化后代码:

import requests
from concurrent.futures import ThreadPoolExecutor

def load_url(url, page, base_parameters):
    # 创建参数副本,避免线程间冲突
    parameters = base_parameters.copy()
    parameters["page"] = page
    try:
        response = requests.get(url=url, params=parameters)
        response.raise_for_status()  # 主动抛出HTTP错误
        return response.json()  # 返回解析后的业务数据
    except Exception as e:
        print(f"页码{page}请求失败: {str(e)}")
        return []

def get_response(total_page_number, url, parameters):
    items = []
    # 并发数取请求数和合理上限的最小值,比如20
    max_workers = min(total_page_number, 20)
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        futures = [
            executor.submit(load_url, url, page, parameters)
            for page in range(1, total_page_number + 1)
        ]
        for future in futures:
            try:
                page_data = future.result()
                items.extend(page_data)
            except Exception as exc:
                print(f"结果处理失败: {str(exc)}")
    return items

更优的异步实现(适配FastAPI)

FastAPI基于asyncio,用异步HTTP客户端aiohttp可彻底避免线程开销,更高效处理IO等待:

import aiohttp
import asyncio

async def fetch_page(session, url, page, base_parameters):
    parameters = base_parameters.copy()
    parameters["page"] = page
    try:
        async with session.get(url, params=parameters) as response:
            response.raise_for_status()
            return await response.json()
    except Exception as e:
        print(f"页码{page}请求失败: {str(e)}")
        return []

async def get_async_response(total_page_number, url, parameters):
    items = []
    # 用信号量控制并发数,避免触发API限流
    semaphore = asyncio.Semaphore(min(total_page_number, 20))
    
    async def bounded_fetch(page):
        async with semaphore:
            async with aiohttp.ClientSession() as session:
                return await fetch_page(session, url, page, parameters)
    
    tasks = [bounded_fetch(page) for page in range(1, total_page_number + 1)]
    results = await asyncio.gather(*tasks)
    
    for page_data in results:
        items.extend(page_data)
    return items

# FastAPI端点示例
from fastapi import FastAPI

app = FastAPI()

@app.get("/fetch-all")
async def fetch_all():
    total_pages = 20
    api_url = "https://external-api.example.com/data"
    base_params = {"per_page": 100}
    all_data = await get_async_response(total_pages, api_url, base_params)
    return {"data": all_data}

额外优化建议

  • 适配API限流规则:先确认外部API的请求频率限制(比如每秒5次),调整并发数避免被限流。
  • 复用HTTP连接:requests用Session对象、aiohttp用ClientSession,复用连接减少TCP握手开销。
  • 流式响应(可选):如果业务允许,用FastAPI的流式响应,获取一页数据就返回一页,不用等全部请求完成。
  • 添加重试机制:对失败请求加指数退避重试,应对临时网络波动。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 22:46:34