异步服务、超时与重试
将 Promise 的经验迁移到 asyncio,并为 AI 服务调用补上超时和重试边界。
学习目标:为异步调用建立时间和取消边界
本节结束时,你能把 JavaScript/TypeScript 的 Promise、await 和 AbortSignal 迁移到 Python asyncio,并解释 coroutine、Task、取消、timeout、retry 和并发预算之间的关系。你会写一个只在暂时性故障上重试的异步调用,为每次等待设置上限,为所有在途任务设置 semaphore,并用可观察的输入、输出和调用次数验证它没有泄漏任务或无限等待。
异步代码的难点不是关键字,而是时间。AI 服务请求可能先等待连接,再等待模型生成,再等待响应体;客户端也可能在中途断开。每个等待点都要有预算,任务被取消时要释放连接和队列位置,重试时要知道服务端是否可能已经执行过操作。把这些边界写出来,才有可运维的推理服务。
从 JS/TS 迁移的心智模型:await 是交出控制权,不是开启线程
在 JS/TS 中,await fetch(...) 会在 promise 未完成时让 event loop 处理其他工作;Python 的 await 也会在 awaitable 未完成时把控制交回事件循环,但它不会把同步函数自动变成异步函数。调用 async def 得到的是 coroutine 对象,直接丢弃它会出现 coroutine was never awaited;调用 asyncio.create_task 才是把它登记为可独立调度的 Task。
const response = await fetch(url, {
signal: AbortSignal.timeout(3000),
});
return await response.json(); async with httpx.AsyncClient(timeout=3) as client:
response = await client.get(url)
response.raise_for_status()
return response.json() 1. coroutine、Task 和 await 的实际生命周期
可以直接 await call_model(text),这样当前任务负责整个调用;也可以先 task = asyncio.create_task(call_model(text)),让事件循环在等待期间运行别的任务。后者带来更高并发,也带来更多责任:成功时读取结果,失败时读取异常,取消时调用 cancel() 并等待 task 收尾。asyncio.gather 负责收集一组结果,但不会替你制定业务上的错误隔离策略。
import asyncio
async def one_request(text: str):
async with asyncio.timeout(3):
return await call_model(text)
async def run_one(text: str):
task = asyncio.create_task(one_request(text), name="model-call")
try:
return await task
except asyncio.CancelledError:
task.cancel()
await asyncio.gather(task, return_exceptions=True)
raise
Python 3.11 的 asyncio.timeout 是异步上下文管理器,离开上下文时会把当前任务的超时取消转换为 TimeoutError;它不能阻止底层系统调用永远不响应,底层客户端仍应有自己的连接和读取 timeout。外层总预算和内层网络预算要相加后检查,不能每次 retry 都重新拿到一份无限时间。
2. timeout:每次等待和整个请求都要有预算
内层 HTTP timeout 用来区分连接慢、读取慢和服务处理慢;外层 asyncio.timeout(total_seconds) 防止重试把用户请求拖到几十秒。若一次调用包含 3 次尝试,每次都设 3 秒,理论上总时间可能超过 9 秒,再加退避;因此可用 deadline 计算剩余时间,或把总 timeout 包住整个 retry 循环。超时错误应带上阶段和 attempt,方便日志定位。
import asyncio
import time
async def with_total_budget(operation, total_seconds: float = 8):
deadline = time.monotonic() + total_seconds
for attempt in range(3):
remaining = deadline - time.monotonic()
if remaining <= 0:
raise TimeoutError("total model-call budget exhausted")
try:
async with asyncio.timeout(remaining):
return await operation()
except TimeoutError:
if attempt == 2:
raise
await asyncio.sleep(min(0.2 * (2 ** attempt), max(0, deadline - time.monotonic())))
注意 asyncio.wait_for(operation(), timeout=...) 会取消被包住的协程并等待它结束,所以实际耗时可能略超过数字;如果 operation 在 finally 中做清理,外层要为清理留余量。timeout 不是吞错:最终要让调用者知道是资源预算耗尽,而不是返回一个伪造的低置信度预测。
3. 取消:正常控制流必须能穿过 retry
客户端断开、服务关闭、任务组失败和部署滚动都可能取消当前 Task。asyncio.CancelledError 不应被普通 except Exception 变成一次可重试的上游失败;清理放在 finally,捕获取消时只做必要的记录和释放,然后重新抛出。持有 semaphore 时,async with 保证退出时归还许可;持有 HTTP response 或文件时也要使用异步上下文管理器。
如果要取消一组任务,保存 task 引用,逐个 cancel(),再用 gather(..., return_exceptions=True) 等待它们完成。否则取消只停了父任务,子任务仍可能继续调用模型服务。对批量 embedding 可以在单项失败后保留其他结果,对一个事务型写操作则通常要整体取消或依靠服务端幂等键恢复。
4. retry:只重试可恢复且安全的操作
重试条件要由错误语义决定。连接超时、暂时断开、502/503/504 可能恢复;401、403、422 和响应 schema 错误通常不能靠重试修复。GET 查询一般更容易幂等,POST 创建任务只有在服务支持幂等键或明确语义时才安全。退避应有上限,加入 jitter 避免所有客户端同时再次请求;最后一次必须保留原始异常和 attempt 信息。
import asyncio
import random
async def retry_timeout(operation, attempts: int = 3):
for attempt in range(attempts):
try:
return await operation()
except asyncio.CancelledError:
raise
except TimeoutError:
if attempt == attempts - 1:
raise
base = min(2.0, 0.2 * (2 ** attempt))
await asyncio.sleep(base + random.uniform(0, 0.05))
真实服务还要处理 Retry-After、429 配额和请求体重放成本。重试计数要进入日志和指标,不能默默把一次请求变成三次 token 消耗。若供应商返回 200 但本地解析失败,应该修复契约或隔离响应,不应把坏数据重试到更大。
5. 并发预算与背压
asyncio.gather(*(predict(x) for x in items)) 会立即创建很多任务,输入来自用户或队列时可能耗尽内存,也可能超过模型服务并发限制。asyncio.Semaphore(n) 限制同时进入调用区的任务数;有限 asyncio.Queue 则让生产者在消费者跟不上时等待。并发预算应与连接池、GPU batch、token 预算、429 速率和下游数据库容量一起测量。
import asyncio
async def bounded_predict(texts: list[str], limit: int = 4):
semaphore = asyncio.Semaphore(limit)
async def one(text: str):
async with semaphore:
try:
value = await retry_timeout(
lambda: asyncio.wait_for(call_model(text), timeout=3),
)
return {"text_id": hash(text), "ok": True, "value": value}
except asyncio.CancelledError:
raise
except Exception as exc:
return {"text_id": hash(text), "ok": False, "error": type(exc).__name__}
return await asyncio.gather(*(one(text) for text in texts))
这里的并发上限限制的是进入模型调用的任务,不代表每个任务都只占一个 token 或一个连接;批量 embedding 仍需要按 token 总量切分。hash(text) 只作示意,生产代码应使用稳定且不暴露原文的样本 ID。结果中的 ok、错误类型和输入 ID 让部分成功可审核,避免一个失败覆盖整批诊断。
运行验证:用 fake operation 观察输出和调用次数
先准备一个 fake operation:前两次抛出 TimeoutError,第三次返回 {"label": "ok"}。运行 retry_timeout 的输入是空参数和 attempts=3,预期输出是该字典、调用次数为 3;若三次都超时,预期最后一次 TimeoutError 被抛出、调用次数为 3。再把任务取消,预期调用次数停止增加,取消传播到外层而不是变成普通失败结果。
对 bounded_predict,用 12 个输入和一个计数器验证 max_in_flight 不超过 4,并检查返回结果数量仍为 12。运行测试时输出 attempts=3, max_in_flight=4, ok=11, failed=1 这样的统计;这些结果比只看到“协程完成”更能证明 timeout、retry 和并发预算真的生效。
常见错误、排错与调试路径
看到 coroutine was never awaited,检查调用点是否遗漏 await 或创建 task 后没有保存;看到所有请求一起变慢,搜索 requests.get、time.sleep 和 CPU 大循环是否进入 async 函数;看到 timeout 之后仍有请求发出,检查 retry 是否捕获了取消或没有等待被取消的 task;看到服务收到三倍流量,检查每层是否重复重试。
调试异步挂起时,打开 asyncio debug,给 task 设置名称,记录阶段、attempt、deadline、队列长度和 semaphore 在途数。若出现 429,先降低并发和 token batch,再看退避;若 p95 延迟只在重试后上升,分别记录首次失败和最终失败。不要在日志里打印完整 prompt,使用样本 ID、长度、哈希和模型版本即可。
练习:实现带总预算的模型调用
任务是设计 predict_with_budget(text):总预算 8 秒,最多两次重试,只对 timeout 重试,每次退避;客户端取消时立即释放 semaphore 和 HTTP 资源。说明一个非幂等写操作为什么不能直接套用同一逻辑,并为“前两次超时后成功”“最终超时”“取消”“并发超过预算”写测试。
为模型服务调用加上有限重试
写一个最多重试两次的异步调用:每次有三秒超时,只对 TimeoutError 重试,并在最后一次保留原异常;为调用增加一个 semaphore。
给我一点提示
用 for attempt in range(3),最后一次失败时直接 raise;不要捕获 CancelledError;外层再用总 deadline 限制整个循环。
查看参考答案
async def predict(text, semaphore):
async with semaphore:
for attempt in range(3):
try:
return await asyncio.wait_for(call_model(text), timeout=3)
except asyncio.CancelledError:
raise
except TimeoutError:
if attempt == 2:
raise
await asyncio.sleep(0.2 * (attempt + 1)) 本节结论
用 fake operation 验证“前两次超时、第三次成功”和“第三次超时仍抛异常”两种路径,再验证取消不会变成第四次调用。若总预算是 8 秒,还要把退避时间算入 deadline;非幂等写操作必须使用幂等键或由服务端提供去重语义。
与 AI 数据管线、模型和推理服务连接
异步适合同时等待多个 HTTP 推理、embedding、检索或数据库操作,但不等于可以无限扩大 batch。每个模型服务都有并发、QPS、token、连接和 GPU 预算;把这些预算写成 semaphore、队列大小和总 timeout,并记录在 job 日志中,才能从结果解释为什么样本被重试、取消或隔离。
在线推理还要把客户端断开与模型拒答分开:前者应取消下游调用并释放资源,后者是正常业务结果;超时和 503 应进入可重试路径,schema 错误应进入供应商协议告警。离线数据管线可以保留部分成功并输出失败清单,在线事务则必须保证重试不会重复扣费、重复写入或重复创建任务。
小结
Python asyncio 的可靠用法是:真正 await 异步操作,给每个等待和整个请求设 timeout,让取消穿过 retry 并释放资源,只在可恢复且安全的条件下重试,用 semaphore 或队列限制并发。运行验证要观察输出、调用次数、在途峰值和取消结果;这些边界共同决定 AI 数据管线和推理服务是否可控。
延伸阅读
先完成本节练习,再用这些资料查阅完整 API 和真实项目组织方式。
阶段共 8 节课,按顺序完成更容易建立完整的迁移模型。