线上采集服务大概每两天挂一次,挂的时候不报错,进程还在,日志停在某一行不动,端口还监听着,但活不干。重启就好,过两个小时再来一遍。
排查过程比想象中久,因为asyncio.wait_for这个函数名太容易让人放心了。它确实会超时,只是超时不等于任务停下来。下面几个用例是把线上代码抽出来之后跑出来的复现,环境 Python 3.12.10 / Windows 11,全部用asyncio.sleep模拟下游,不依赖网络,跑出来的数字和文章里写的一致。完整代码在文末。
2026-09-18 04:12:07 INFO [xianlin] batch 7731 start, 240 tasks 2026-09-18 04:12:08 INFO [xianlin] task 12 timeout after 5.0s 2026-09-18 06:31:44 INFO [xianlin] batch 7731 start, 240 tasks04:12 那一行之后,日志空白了两个多小时。注意第二行看着很健康——它确实打印了超时,然后什么事都没发生。
更麻烦的一点是进程状态完全正常。top看 CPU 不高,内存平稳,netstat显示端口在听,健康检查接口返回 200——因为事件循环还在转,只是里面塞了一堆永远不会结束的任务。监控上唯一异常的是"这批任务的处理条数",但它掉得慢,告警阈值要两小时才碰到,那时候已经晚了。
先确认你中招了没有
在往下看之前,有个五分钟的自检能判断你的服务有没有同样的问题。在事件循环里加一段,看看有没有任务处于"收到取消但没退出"的状态:
importasyncioasyncdefwatchdog()->None:"""每 10 秒扫一遍, 抓那些取消了却还活着的任务。"""whileTrue:awaitasyncio.sleep(10)alive=[tfortinasyncio.all_tasks()ifnott.done()]# 挂着超过 60 秒的任务, 正常业务里不该存在stale=[tfortinaliveift.get_name().startswith("fetch-")]iflen(stale)>50:print(f"[warn] 疑似泄漏任务{len(stale)}个, 示例:{stale[0].get_name()}")判定思路是看任务名和存活时长,而不是看错误日志——这种故障的最大特征就是不产生错误日志。如果你的服务在某些时段任务数只增不减,基本可以确认下面某一条踩中了。
超时被协程自己吞掉了
最容易出事的一种。有人在协程里写了"取消时优雅收尾",本意是好的:
asyncdeffetch()->str:try:awaitasyncio.sleep(3.0)return"ok"exceptasyncio.CancelledError:# 本意是取消时把中间状态落库, 结果把外层超时一起废掉了return"finished-anyway"实测这一段的行为是:
期望: 1.00s 抛 TimeoutError 实测: 1.01s, 正常返回 'finished-anyway'wait_for在 1 秒时确实动手了,它给协程发了取消信号。协程把CancelledError接住,然后返回了一个正常值。于是wait_for认为任务完成了,把"finished-anyway"当成结果交给调用方。
调用方等满 1 秒,拿到一个来路不明的返回值,没有任何异常,日志里也不会出现 TimeoutError。上游看到的是成功。这是最难受的故障形态——不是慢,是假成功。
判断标准很简单:except asyncio.CancelledError后面如果跟的不是raise,这个地方就有问题。想收尾可以,收完必须重新抛出:
exceptasyncio.CancelledError:awaitself.flush_state()# 收尾随便做raise# 这一行不能少shield 会让取消彻底送不到
顺着上一条往下查,还能发现一个更隐蔽的情况。有些代码怕取消影响到下游,会套一层asyncio.shield:
awaitasyncio.wait_for(asyncio.shield(slow()),timeout=0.5)这种写法看着稳妥,实际把整条取消链路切断了。实测了一下,在下游放一个探针,看它的except CancelledError到底有没有被触发:
写法A: 0.51s 抛 TimeoutError, 此刻下游日志 空 再等 1.8s 后下游日志: ['downstream-finished']关键在最后一行。wait_for在 0.51 秒抛了 TimeoutError,调用方以为超时生效了;但下游从头到尾没收到取消信号,既没有downstream-cancelled,也没有任何异常,它是正常跑完的,日志打的是downstream-finished。
也就是说,上一条讲的"收尾逻辑必须重新raise",在套了 shield 之后连触发的机会都没有——因为下游压根不知道自己在被超时。
shield的设计目的是"别让取消打断我正在进行的关键操作",这个目的一点问题都没有。问题在于很多人把它当成了"加个保护更安全",顺手套上,然后指望外面的wait_for还能叫停里面。这两件事互斥。
如果确实需要 shield,正确做法是给下游自己的超时,让它在内部主动停:
asyncdefslow_with_own_timeout()->str:try:asyncwithasyncio.timeout(0.5):# 下游自己知道什么时候该停awaitasyncio.sleep(2.0)return"late-result"exceptasyncio.TimeoutError:log("downstream-self-timeout")# 收尾在这里做, 不依赖外部取消raise实测这一版的下游日志是['downstream-self-timeout'],收尾正常执行。区别在于:外层的超时是给调用方看的,下游的超时才是给下游用的。套了 shield 之后,前者管不到后者。
单次调用的超时,管不了整个任务
第二类问题跟取消无关,纯粹是数学问题。超时写在重试循环里面:
asyncdefwith_retry()->str:last=Nonefor_inrange(4):try:returnawaitasyncio.wait_for(call_once(),timeout=1.0)exceptasyncio.TimeoutErrorasexc:last=excraiselast每次调用最多 1 秒,看起来整个任务最多……4 秒。实测:
期望: 整个任务最多 1.00s 实测: 4.03s, 重试了 4 次1 秒的超时乘上重试次数,变成 4 秒。上游网关给这个接口的预算是 2 秒,超过就断连。所以第 3、4 次重试的结果永远没人接收,纯属浪费——而且这 4 秒里 worker 一直被占着。
超时要分两层看:单次调用的上限,和整个函数的预算。前者防下游单次抖动,后者防重试累加。只写前者,总耗时就是"上限 × 重试次数"。
取消信号本身是能穿透的
排查过程中一度怀疑except Exception会顺手把CancelledError一起吃掉。实测了一下,在 Python 3.12 上不会:
取消传递链路: inner-cancelled -> outer-saw-cancel -> caller-saw-cancelCancelledError从 Python 3.8 起继承自BaseException,不走except Exception这条路径。这一条可以让人放心:不用怕except Exception吃掉取消,要怕的是显式写了except CancelledError又没raise的地方。
executor 里的任务,取消不掉
真正让 worker 被占死的是这一类。有些老 SDK 只有同步接口,只能丢进run_in_executor。给它 1 秒超时,它确实 1 秒就返回了:
任务A: 1.01s 时超时返回, 但 worker 线程还在跑 blocking_io 任务B: 本身只要 0.10s, 实际等了 2.09s任务 A 超时返回了,看着没问题。但底层那个线程不认 asyncio 的取消标记,它老老实实把 3 秒的time.sleep跑完。用max_workers=1复现了池子只有 1 个 worker 的情况:任务 B 本身只要 0.1 秒,实际等了 2.09 秒,纯粹在排队。
线上 worker 数量是固定的。每个超时返回的请求都还占着一个 worker,超时越多占得越多,剩下的 worker 越少,新请求排队越久,超时更多。两个小时的雪崩就是这么滚起来的。
要限住这种任务,靠 asyncio 的超时没用,得从线程池本身下手:给池子设一个够小的max_workers让压力在队列层面暴露出来,或者在同步函数内部自己做超时控制(比如 httpx 的同步客户端能设 timeout)。把不可取消的任务塞进一个无限大的默认线程池,等于把问题藏起来,等它攒够了再一起爆。
三层预算的写法
把上面几个坑合成一个可复用的模式,这段可以直接抄:
asyncdeffetch_with_budget(url:str,total_budget:float,delay:float=0.4)->str:deadline=asyncio.get_running_loop().time()+total_budgetasyncdefonce()->str:remaining=deadline-asyncio.get_running_loop().time()ifremaining<=0:raiseasyncio.TimeoutError("总预算已耗尽")# 单次上限取"剩余预算"和"单次上限"的较小值,# 否则最后一次重试必然打穿总时长asyncwithasyncio.timeout(min(1.0,remaining)):awaitasyncio.sleep(delay)returnf"body-of-{url}"last=Nonefor_inrange(3):try:returnawaitonce()exceptasyncio.TimeoutErrorasexc:last=excexceptasyncio.CancelledError:raise# 取消原样抛出raiselast跑出来的结果:
场景A 下游 0.15s / 预算 0.6s: 0.16s 成功 场景B 下游 1.2s / 预算 0.5s: 0.50s 抛出 TimeoutError(总预算已耗尽), 只尝试了 3 次场景 B 是关键:下游每次要 1.2 秒,预算给 0.5 秒,函数在 0.5 秒整放弃。这里没有"1 秒超时 × 3 次重试"的累加,因为每次重试前都会先看还剩多少预算,第二次进来remaining已经是负数,直接抛。
三层各管一件事:
| 层次 | 管什么 | 典型值 |
|---|---|---|
| 单次调用上限 | 挡住下游偶发抖动 | 1-5 秒 |
| 整个函数预算 | 挡住重试累加 | 上游网关超时 × 0.6 |
| 可穿透的取消 | 保证外层能真的叫停 | 不吞CancelledError |
总预算那个 0.6 是经验值。上游网关通常给 10 秒,那这个函数内部就别超过 6 秒,留出序列化、写日志、返回响应的时间。把上游的断连时间当成自己的预算,是这几条里最容易忘的一条。
顺带说一句版本问题:asyncio.timeout()这个上下文管理器是 Python 3.11 才有的,3.10 及以下只能用asyncio.wait_for,写的时候注意项目和运行时的版本对不对得上。
要给别人的协程加超时,用包装器
上面那个fetch_with_budget需要你改自己的函数体。但需求常常是反过来的:想给一个已经写好、改动不了的第三方协程套上这套预算,比如 SDK 里现成的client.fetch()。
这种时候别去动原函数,写个包装器:
importasynciofromcollections.abcimportAwaitable,CallablefromtypingimportTypeVar T=TypeVar("T")asyncdefwith_budget(factory:Callable[[],Awaitable[T]],total_budget:float,attempts:int=3,label:str="call",)->T:""" 给任意协程工厂套三层预算。传工厂而不是协程对象, 因为协程对象只能 await 一次, 重试需要每次重新创建。 注意参数是 factory 不是 coro —— 这是这个函数最容易用错的地方。 """deadline=asyncio.get_running_loop().time()+total_budget last:BaseException|None=Noneforiinrange(1,attempts+1):remaining=deadline-asyncio.get_running_loop().time()ifremaining<=0:raiseasyncio.TimeoutError(f"{label}总预算{total_budget}s 耗尽")try:asyncwithasyncio.timeout(remaining):returnawaitfactory()exceptasyncio.CancelledError:raise# 取消原样抛, 不参与重试except(asyncio.TimeoutError,TimeoutError)asexc:last=exc log(f"{label}第{i}次超时, 剩余预算{remaining:.2f}s")raiselast# type: ignore[misc]用起来是这样,原来的协程一行都不用改:
body=awaitwith_budget(lambda:client.fetch("https://example.com/api"),total_budget=2.0,label="fetch-api",)有两个细节值得单独说。参数必须是工厂(lambda: client.fetch(...))而不是协程对象(client.fetch(...)),因为协程对象只能被 await 一次,重试的时候第二次 await 会直接报RuntimeError: cannot reuse already awaited coroutine,这个错还挺容易误判成 SDK 的问题。另外except CancelledError: raise必须放在except TimeoutError前面,否则取消有可能被后面的分支吃进去。
不过要提醒一句:包装器只能管住"什么时候放弃等待",管不住"下游什么时候真的停"。它和前面 executor 那个问题是同一类——超时是调用方的决定,不是被调用方的行为。要真正止损,还是得让下游自己能超时。
一个例外
except CancelledError后面不raise一定是错的,但有极少数场景你确实需要拦住取消。比如正在写数据库事务,中途取消会留下半提交状态。这种时候正确做法是拦住、把清理做完、再抛出去:
exceptasyncio.CancelledError:asyncwithself._transaction_lock:awaitself.rollback()# 清理必须做完raise# 然后照样抛, 不能改成 return区别在于:raise之后外层知道这个任务被取消了,return之后外层以为它成功了。前者是延迟响应取消,后者是伪造成功。
现在可以打开你项目里的asyncio.wait_for,按顺序看四件事:括号里面的协程有没有except CancelledError且没raise;有没有被asyncio.shield包住;这个超时是不是写在重试循环内部;被包起来的协程里有没有run_in_executor。这四个问题在这种超时失效的故障里能覆盖绝大多数情况。查完把结论记到代码注释里,比下次故障再翻一遍强。
附:完整复现代码
存成timeout_demo.py直接python timeout_demo.py就能跑,无第三方依赖:
#!/usr/bin/env python# -*- coding: utf-8 -*-""" 演示 asyncio.wait_for 在几种写法下"超时失效"的完整可运行代码。 运行环境: Python 3.12.10 / Windows 11 """from__future__importannotationsimportasyncioimportfunctoolsimportsysimporttimefromconcurrent.futuresimportThreadPoolExecutor DOWNSTREAM_DELAY=3.0# 下游实际要花 3 秒PER_CALL_TIMEOUT=1.0# 我们希望单次调用最多等 1 秒asyncdefcase1_swallowed_cancel()->None:print("用例 1: 协程吞掉 CancelledError, 超时静默变成假成功")asyncdefslow()->str:try:awaitasyncio.sleep(3.0)return"ok"exceptasyncio.CancelledError:# 本意是"取消时优雅收尾", 结果把外层超时一起废掉了return"finished-anyway"t0=time.perf_counter()try:result=awaitasyncio.wait_for(slow(),timeout=1.0)outcome=f"正常返回{result!r}"exceptasyncio.TimeoutError:outcome="抛出了 TimeoutError"print(f" 期望: 1.00s 抛 TimeoutError")print(f" 实测:{time.perf_counter()-t0:.2f}s,{outcome}")asyncdefcase2_retry_budget()->None:print("用例 2: 超时放在重试里面, 总耗时失控")attempts=0asyncdefcall_once()->str:nonlocalattempts attempts+=1awaitasyncio.sleep(PER_CALL_TIMEOUT+0.2)# 每次都刚好超时return"ok"asyncdefwith_retry()->str:last:Exception|None=Nonefor_inrange(4):try:returnawaitasyncio.wait_for(call_once(),timeout=PER_CALL_TIMEOUT)exceptasyncio.TimeoutErrorasexc:last=excraiselast# type: ignore[misc]t0=time.perf_counter()try:awaitwith_retry()exceptasyncio.TimeoutError:passprint(f" 期望: 整个任务最多 1.00s")print(f" 实测:{time.perf_counter()-t0:.2f}s, 重试了{attempts}次")asyncdefcase3_cancel_propagation()->None:print("用例 3: except Exception 是否会吞掉取消信号")stage:list[str]=[]asyncdefinner()->None:try:awaitasyncio.sleep(10)exceptasyncio.CancelledError:stage.append("inner-cancelled")raiseasyncdefouter()->None:try:awaitinner()exceptasyncio.CancelledError:stage.append("outer-saw-cancel")raisetask=asyncio.create_task(outer())awaitasyncio.sleep(0.05)task.cancel()try:awaittaskexceptasyncio.CancelledError:stage.append("caller-saw-cancel")print(f" 取消传递链路:{' -> '.join(stage)}")defblocking_io(seconds:float)->str:"""同步阻塞调用, 模拟某个只能用同步库的 SDK。"""time.sleep(seconds)return"done"asyncdefcase4_executor_not_cancellable()->None:print("用例 4: executor 里的阻塞任务取消不掉, 会连累后面的任务")loop=asyncio.get_running_loop()# 只给 1 个线程: 线上连接池、线程池被打满就是这个效果withThreadPoolExecutor(max_workers=1)aspool:t0=time.perf_counter()fut_a=loop.run_in_executor(pool,blocking_io,3.0)try:awaitasyncio.wait_for(fut_a,timeout=1.0)exceptasyncio.TimeoutError:print(f" 任务A:{time.perf_counter()-t0:.2f}s 时超时返回, "f"但 worker 线程还在跑 blocking_io")t1=time.perf_counter()awaitloop.run_in_executor(pool,blocking_io,0.1)print(f" 任务B: 本身只要 0.10s, 实际等了{time.perf_counter()-t1:.2f}s")asyncdefcase6_shield_breaks_timeout()->None:print("用例 6: 用 shield 保护下游, 取消送不到, 收尾也不会触发")logs:list[str]=[]asyncdefslow()->str:try:awaitasyncio.sleep(2.0)logs.append("downstream-finished")return"late-result"exceptasyncio.CancelledError:logs.append("downstream-cancelled")raiset0=time.perf_counter()try:awaitasyncio.wait_for(asyncio.shield(slow()),timeout=0.5)exceptasyncio.TimeoutError:print(f" 写法A:{time.perf_counter()-t0:.2f}s 抛 TimeoutError, "f"此刻下游日志{logsor'空'}")awaitasyncio.sleep(1.8)print(f" 再等 1.8s 后下游日志:{logs}")asyncdeffetch_with_budget(url:str,total_budget:float,delay:float=0.4)->str:"""外层给总预算, 内层给单次上限, 取消必须能穿透。"""deadline=asyncio.get_running_loop().time()+total_budgetasyncdefonce()->str:remaining=deadline-asyncio.get_running_loop().time()ifremaining<=0:raiseasyncio.TimeoutError("总预算已耗尽")asyncwithasyncio.timeout(min(PER_CALL_TIMEOUT,remaining)):awaitasyncio.sleep(delay)returnf"body-of-{url}"last:Exception|None=Noneattempts=0for_inrange(3):attempts+=1try:body=awaitonce()fetch_with_budget.last_attempts=attemptsreturnbodyexceptasyncio.TimeoutErrorasexc:last=excexceptasyncio.CancelledError:raisefetch_with_budget.last_attempts=attemptsraiselast# type: ignore[misc]fetch_with_budget.last_attempts=0asyncdefcase5_correct_pattern()->None:print("用例 5: 带总预算的正确写法")t0=time.perf_counter()body=awaitfetch_with_budget("https://example.com/api",total_budget=0.6,delay=0.15)print(f" 场景A 下游 0.15s / 预算 0.6s: "f"{time.perf_counter()-t0:.2f}s 成功,{body!r}")t0=time.perf_counter()try:awaitfetch_with_budget("https://example.com/down",total_budget=0.5,delay=1.2)print(" 场景B: 居然成功了, 说明预算没生效")exceptasyncio.TimeoutErrorasexc:print(f" 场景B 下游 1.2s / 预算 0.5s:{time.perf_counter()-t0:.2f}s "f"抛出 TimeoutError({exc}), 只尝试了{fetch_with_budget.last_attempts}次")asyncdefmain()->None:# 让被丢弃的任务不要往控制台喷 "Task exception was never retrieved"asyncio.get_running_loop().set_exception_handler(lambdaloop,ctx:None)print(f"Python{sys.version.split()[0]}单次调用超时设为{PER_CALL_TIMEOUT}s\n")awaitcase1_swallowed_cancel()awaitcase2_retry_budget()awaitcase3_cancel_propagation()awaitcase4_executor_not_cancellable()awaitcase6_shield_breaks_timeout()awaitcase5_correct_pattern()if__name__=="__main__":asyncio.run(main())跑完的输出应该和上面每段引用的实测结果一致。如果某一段对不上,先检查 Python 版本——3.11 以下没有asyncio.timeout(),需要把用例 5 换回asyncio.wait_for。
下一步建议在你自己项目里做一件事:把grep -rn "wait_for" --include=*.py的结果过一遍,凡是超时参数写在重试循环里面的,按本文的三层预算改成总预算版本。这一步改动通常不超过二十行,但能挡掉大部分"日志突然不动"的故障。