FastAPI中如何处理可能永不回调的异步外部API?
问题:FastAPI调用异步回调API的超时处理方案
背景与问题
我用FastAPI构建API服务,需调用一款名为ex_api的异步外部API,调用流程为:
- 我的服务向ex_api的创建任务端点发送请求,参数包含
callback_url; - ex_api的创建任务端点立即返回含任务ID的响应;
- ex_api后台静默处理任务;
- 处理完成后,ex_api向
callback_url发送携带任务结果的请求; - 我的服务收到请求后返回响应,流程结束。
当前问题是ex_api稳定性不足,有时会卡在第3步无法触发回调,我无法感知该情况。希望设置有限等待时长,超时未收到回调则执行其他方案,但不知如何实现,且无权修改ex_api的任何内容。
我的疑问:
- 有哪些相关的最佳实践或经验可以分享?
- 我的超时等待方案是否可实现?若可以,该如何操作?
问题代码示例
from typing import Any import fastapi import httpx from pydantic import BaseModel THIS_SERVER_URL = "https://fake.this.server.public.url" CALLBACK_API = "/callback/plus_one" CALLBACK_URL = f"{THIS_SERVER_URL}{CALLBACK_API}" API_PLUS_ONE = "https://fake.api.url/api/plus_one" app = fastapi.FastAPI() class MyQuery(BaseModel): x: int class MyResponse(BaseModel): message: str result: int @app.post("/api/plus_one_mul_two", response_model=MyResponse) async def post_plus_one_mul_two(params: MyQuery): async with httpx.AsyncClient() as c: response = await c.post( API_PLUS_ONE, json={ "callback_url": CALLBACK_URL, "number": params.x, }, ) response.raise_for_status() result: dict[str, Any] = response.json() return {"message": result["message"]} class CbQuery(BaseModel): message: str number: int class CbResponse(BaseModel): message: str async def fake_upload_to_oss(number: int): pass # FIXME: May never receive callback! @app.post(CALLBACK_API, response_model=CbResponse) async def callback_plus_one_mul_two(params: CbQuery): if params.message == "ok": await fake_upload_to_oss(params.number * 2) else: raise Exception("failed") return {"message": "ok"}
解答
一、相关最佳实践
- 任务状态追踪:每次请求ex_api时,将任务ID、请求参数、创建时间、状态(待处理/已完成/超时/失败)存入数据库或分布式缓存(如Redis),方便后续校验、重试和状态查询。
- 主动轮询兜底:若回调超时,优先调用ex_api的任务查询接口(如果存在)获取状态;若无查询接口,可对超时任务进行重试(需保证请求幂等,比如给每个请求加唯一标识,避免重复处理)。
- 幂等性设计:回调接口和重试逻辑必须保证幂等,比如用任务ID作为唯一键,重复收到相同任务的回调时直接返回成功,不重复执行后续业务操作。
- 合理设置超时阈值:根据ex_api的平均处理时间设置超时时间(比如平均处理10秒,设15-20秒超时),避免误判;同时设置最大重试次数,防止无限重试消耗资源。
- 告警机制:统计超时任务占比,当超过预设阈值时触发告警(如邮件、企业微信通知),及时发现ex_api的稳定性问题。
二、超时等待方案的实现
你的需求完全可以实现,核心思路是发起ex_api请求后,启动异步超时等待任务,通过分布式缓存共享任务状态,超时则执行降级逻辑。以下是具体实现代码:
1. 依赖引入与基础配置
from typing import Any, Optional import fastapi import httpx import asyncio from pydantic import BaseModel import redis.asyncio as redis THIS_SERVER_URL = "https://fake.this.server.public.url" CALLBACK_API = "/callback/plus_one" API_PLUS_ONE = "https://fake.api.url/api/plus_one" # 超时时间(单位:秒) TIMEOUT_SECONDS = 20 # Redis连接配置(多实例部署必须用分布式缓存) REDIS_URL = "redis://localhost:6379/0" app = fastapi.FastAPI() # 初始化Redis客户端 redis_client = redis.from_url(REDIS_URL)
2. 修改主接口逻辑
发起ex_api请求后,等待回调结果或超时,超时则执行降级方案:
class MyQuery(BaseModel): x: int class MyResponse(BaseModel): message: str result: Optional[int] = None @app.post("/api/plus_one_mul_two", response_model=MyResponse) async def post_plus_one_mul_two(params: MyQuery): async with httpx.AsyncClient() as c: response = await c.post( API_PLUS_ONE, json={ "number": params.x, }, ) response.raise_for_status() result_data = response.json() task_id = result_data.get("task_id") if not task_id: return {"message": "Failed to get task ID", "result": None} # 构造带task_id的回调URL(若ex_api支持在请求体返回task_id,也可让回调携带) callback_url = f"{THIS_SERVER_URL}{CALLBACK_API}?task_id={task_id}" # 重新发起请求,携带带task_id的回调URL async with httpx.AsyncClient() as c: await c.post( API_PLUS_ONE, json={ "callback_url": callback_url, "number": params.x, }, ) # 初始化任务状态:待处理,超时时间比等待时间多5秒,避免缓存提前过期 await redis_client.setex(f"task:{task_id}", TIMEOUT_SECONDS + 5, "pending") try: # 等待回调结果,超时触发异常 await asyncio.wait_for( wait_for_callback(task_id), timeout=TIMEOUT_SECONDS ) # 获取回调返回的结果 result_number = await redis_client.get(f"task_result:{task_id}") if result_number: final_result = int(result_number) * 2 await fake_upload_to_oss(final_result) # 清理缓存 await redis_client.delete(f"task:{task_id}", f"task_result:{task_id}") return {"message": "Success", "result": final_result} else: return {"message": "Callback received but no result", "result": None} except asyncio.TimeoutError: # 超时降级逻辑:这里用本地模拟加1计算作为示例 await redis_client.setex(f"task:{task_id}", 3600, "timeout") fallback_result = (params.x + 1) * 2 await fake_upload_to_oss(fallback_result) return {"message": "Callback timeout, fallback executed", "result": fallback_result} async def wait_for_callback(task_id: str): # 循环检查任务状态,直到完成或超时 while True: status = await redis_client.get(f"task:{task_id}") if status == b"completed": break await asyncio.sleep(0.5)
3. 修改回调接口逻辑
回调收到结果后,更新任务状态并存储结果:
class CbQuery(BaseModel): message: str number: int class CbResponse(BaseModel): message: str async def fake_upload_to_oss(number: int): pass @app.post(CALLBACK_API, response_model=CbResponse) async def callback_plus_one_mul_two(params: CbQuery, task_id: str = fastapi.Query(...)): if params.message == "ok": # 存储任务结果 await redis_client.setex(f"task_result:{task_id}", 3600, str(params.number)) # 更新任务状态为已完成 await redis_client.set(f"task:{task_id}", "completed") else: await redis_client.set(f"task:{task_id}", "failed") raise fastapi.HTTPException(status_code=500, detail="Task failed") return {"message": "ok"}
关键注意点
- task_id传递:必须保证回调能携带task_id,若ex_api不支持在请求体返回,可在构造
callback_url时将task_id作为查询参数拼接进去。 - 分布式部署适配:多实例部署时不能用内存变量存储状态,必须用Redis等分布式缓存,否则不同实例无法共享任务状态。
- 降级逻辑适配:降级方案需根据业务场景设计,比如本地模拟计算、调用备用API等,确保业务可用性。
内容的提问来源于stack exchange,提问作者Isuxiz Slidder
相关产品推荐
相关产品推荐

