如何用FastAPI处理同时发起的重复请求仅执行一次?
问题描述
我有一个运行在10.11.12.13:8000的FastAPI应用,同时注册了多个使用该地址的worker。客户端可通过注册地址向worker发送请求(每个请求都是唯一的),但所有worker都注册了同一地址,因此所有请求都会发送到10.11.12.13:8000。
现在需要处理的问题是:客户端会通过asyncio.gather同时向多个worker发送同一个唯一请求。如何让我的应用仅处理该请求一次,却给客户端返回仿佛每个worker都独立完成任务的响应?
我曾考虑对短时间内到来的请求进行批处理,或者使用缓存,但每个请求都是唯一的,所以我认为缓存不是好办法。
示例代码
app.py
from fastapi import FastAPI import uvicorn app = FastAPI() work_count = 0 @app.get("/") def handle_workload(): global work_count work_count += 1 return {"message": f"Hello world! {work_count}"} def start_server(): uvicorn.run( app, host="0.0.0.0", port=8000, log_level="debug" ) if __name__ == "__main__": start_server()
client.py
import asyncio import aiohttp import random async def send_request(worker_url: str): async with aiohttp.ClientSession() as session: async with session.get(worker_url) as response: if response.status == 200: return await response.json() # 所有worker都注册了同一个地址 worker_servers = [ ("worker_1", "http://localhost:8000"), ("worker_2", "http://localhost:8000"), ("worker_3", "http://localhost:8000"), ("worker_4", "http://localhost:8000") ] async def main(): worker_urls = [worker[1] for worker in random.sample(worker_servers, 3)] responses = await asyncio.gather(*[send_request(url) for url in worker_urls]) print(responses) if __name__ == "__main__": asyncio.run(main())
实现建议
1. 异步任务共享机制(核心方案)
基于请求唯一标识,维护正在执行的任务缓存,让重复请求复用同一任务结果:
- 用全局字典存储正在处理的异步任务,键为请求的唯一标识(可根据请求参数、业务ID生成),值为
asyncio.Task对象。 - 用异步锁保护任务字典的读写,避免并发冲突。
- 请求到达时先检查缓存:存在则等待任务完成复用结果;不存在则创建新任务,完成后自动清理缓存。
修改后的app.py示例:
from fastapi import FastAPI import uvicorn import asyncio from typing import Dict app = FastAPI() # 存储正在处理的任务:键为请求唯一标识 active_tasks: Dict[str, asyncio.Task] = {} work_count = 0 # 保护任务字典的异步锁 task_lock = asyncio.Lock() async def actual_work(): """实际执行业务逻辑的异步函数""" global work_count await asyncio.sleep(0.1) # 模拟耗时操作 work_count += 1 return {"message": f"Hello world! {work_count}"} @app.get("/") async def handle_workload(): # 实际场景中可根据请求参数/业务ID生成唯一标识 request_key = "unique_business_key" async with task_lock: task = active_tasks.get(request_key) if not task: task = asyncio.create_task(actual_work()) active_tasks[request_key] = task # 任务完成后自动清理缓存 task.add_done_callback(lambda t: asyncio.create_task(_cleanup(request_key))) # 等待任务完成,返回结果 return await task async def _cleanup(request_key: str): async with task_lock: active_tasks.pop(request_key, None) def start_server(): uvicorn.run( app, host="0.0.0.0", port=8000, log_level="debug" ) if __name__ == "__main__": start_server()
2. 客户端侧去重(如果客户端代码可控)
直接在客户端避免发送重复请求,复制结果模拟多worker响应:
- 对要发送的URL列表去重,只发送一次请求。
- 将结果复制对应次数,返回给
asyncio.gather的调用逻辑。
修改后的client.py示例:
import asyncio import aiohttp import random async def send_request(worker_url: str): async with aiohttp.ClientSession() as session: async with session.get(worker_url) as response: if response.status == 200: return await response.json() worker_servers = [ ("worker_1", "http://localhost:8000"), ("worker_2", "http://localhost:8000"), ("worker_3", "http://localhost:8000"), ("worker_4", "http://localhost:8000") ] async def main(): worker_urls = [worker[1] for worker in random.sample(worker_servers, 3)] # 去重后只保留唯一请求地址 unique_url = list(set(worker_urls))[0] # 发送一次请求 result = await send_request(unique_url) # 复制结果模拟多个worker响应 responses = [result] * len(worker_urls) print(responses) if __name__ == "__main__": asyncio.run(main())
3. 短窗口批处理(针对集中到达的重复请求)
设置一个极短的时间窗口,收集同一业务标识的请求,批量返回结果:
- 用定时器或滑动窗口,在窗口内收集相同业务标识的请求。
- 窗口结束后执行一次任务,将结果返回给所有等待的请求。
- 适合请求集中爆发的场景,需定义明确的业务请求标识。
关键注意事项
- 唯一标识必须精准:确保同一业务逻辑的请求对应同一个标识,不同业务请求标识不重复,避免结果复用错误。
- 缓存及时清理:任务成功/失败后都要清理缓存,防止内存泄漏。
- 异常处理:任务执行失败时,要让所有等待的请求都能捕获到错误,避免请求挂起。
内容的提问来源于stack exchange,提问作者nntoan209
相关产品推荐
相关产品推荐

