如何优化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
相关产品推荐
相关产品推荐

