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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 12:44:53