2026/9/24 20:51:24

asyncio 超时设错,我的采集服务每天静默挂两小时

asyncio 超时设错,我的采集服务每天静默挂两小时 线上采集服务大概每两天挂一次挂的时候不报错进程还在日志停在某一行不动端口还监听着但活不干。重启就好过两个小时再来一遍。排查过程比想象中久因为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)returnokexceptasyncio.CancelledError:# 本意是取消时把中间状态落库, 结果把外层超时一起废掉了returnfinished-anyway实测这一段的行为是期望: 1.00s 抛 TimeoutError 实测: 1.01s, 正常返回 finished-anywaywait_for在 1 秒时确实动手了它给协程发了取消信号。协程把CancelledError接住然后返回了一个正常值。于是wait_for认为任务完成了把finished-anyway当成结果交给调用方。调用方等满 1 秒拿到一个来路不明的返回值没有任何异常日志里也不会出现 TimeoutError。上游看到的是成功。这是最难受的故障形态——不是慢是假成功。判断标准很简单except asyncio.CancelledError后面如果跟的不是raise这个地方就有问题。想收尾可以收完必须重新抛出exceptasyncio.CancelledError:awaitself.flush_state()# 收尾随便做raise# 这一行不能少shield 会让取消彻底送不到顺着上一条往下查还能发现一个更隐蔽的情况。有些代码怕取消影响到下游会套一层asyncio.shieldawaitasyncio.wait_for(asyncio.shield(slow()),timeout0.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)returnlate-resultexceptasyncio.TimeoutError:log(downstream-self-timeout)# 收尾在这里做, 不依赖外部取消raise实测这一版的下游日志是[downstream-self-timeout]收尾正常执行。区别在于外层的超时是给调用方看的下游的超时才是给下游用的。套了 shield 之后前者管不到后者。单次调用的超时管不了整个任务第二类问题跟取消无关纯粹是数学问题。超时写在重试循环里面asyncdefwith_retry()-str:lastNonefor_inrange(4):try:returnawaitasyncio.wait_for(call_once(),timeout1.0)exceptasyncio.TimeoutErrorasexc:lastexcraiselast每次调用最多 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_workers1复现了池子只有 1 个 worker 的情况任务 B 本身只要 0.1 秒实际等了 2.09 秒纯粹在排队。线上 worker 数量是固定的。每个超时返回的请求都还占着一个 worker超时越多占得越多剩下的 worker 越少新请求排队越久超时更多。两个小时的雪崩就是这么滚起来的。要限住这种任务靠 asyncio 的超时没用得从线程池本身下手给池子设一个够小的max_workers让压力在队列层面暴露出来或者在同步函数内部自己做超时控制比如 httpx 的同步客户端能设 timeout。把不可取消的任务塞进一个无限大的默认线程池等于把问题藏起来等它攒够了再一起爆。三层预算的写法把上面几个坑合成一个可复用的模式这段可以直接抄asyncdeffetch_with_budget(url:str,total_budget:float,delay:float0.4)-str:deadlineasyncio.get_running_loop().time()total_budgetasyncdefonce()-str:remainingdeadline-asyncio.get_running_loop().time()ifremaining0:raiseasyncio.TimeoutError(总预算已耗尽)# 单次上限取剩余预算和单次上限的较小值,# 否则最后一次重试必然打穿总时长asyncwithasyncio.timeout(min(1.0,remaining)):awaitasyncio.sleep(delay)returnfbody-of-{url}lastNonefor_inrange(3):try:returnawaitonce()exceptasyncio.TimeoutErrorasexc:lastexcexceptasyncio.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 TTypeVar(T)asyncdefwith_budget(factory:Callable[[],Awaitable[T]],total_budget:float,attempts:int3,label:strcall,)-T: 给任意协程工厂套三层预算。传工厂而不是协程对象, 因为协程对象只能 await 一次, 重试需要每次重新创建。 注意参数是 factory 不是 coro —— 这是这个函数最容易用错的地方。 deadlineasyncio.get_running_loop().time()total_budget last:BaseException|NoneNoneforiinrange(1,attempts1):remainingdeadline-asyncio.get_running_loop().time()ifremaining0:raiseasyncio.TimeoutError(f{label}总预算{total_budget}s 耗尽)try:asyncwithasyncio.timeout(remaining):returnawaitfactory()exceptasyncio.CancelledError:raise# 取消原样抛, 不参与重试except(asyncio.TimeoutError,TimeoutError)asexc:lastexc log(f{label}第{i}次超时, 剩余预算{remaining:.2f}s)raiselast# type: ignore[misc]用起来是这样原来的协程一行都不用改bodyawaitwith_budget(lambda:client.fetch(https://example.com/api),total_budget2.0,labelfetch-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就能跑无第三方依赖展开完整代码约 200 行#!/usr/bin/env python# -*- coding: utf-8 -*- 演示 asyncio.wait_for 在几种写法下超时失效的完整可运行代码。 运行环境: Python 3.12.10 / Windows 11 from__future__importannotationsimportasyncioimportfunctoolsimportsysimporttimefromconcurrent.futuresimportThreadPoolExecutor DOWNSTREAM_DELAY3.0# 下游实际要花 3 秒PER_CALL_TIMEOUT1.0# 我们希望单次调用最多等 1 秒asyncdefcase1_swallowed_cancel()-None:print(用例 1: 协程吞掉 CancelledError, 超时静默变成假成功)asyncdefslow()-str:try:awaitasyncio.sleep(3.0)returnokexceptasyncio.CancelledError:# 本意是取消时优雅收尾, 结果把外层超时一起废掉了returnfinished-anywayt0time.perf_counter()try:resultawaitasyncio.wait_for(slow(),timeout1.0)outcomef正常返回{result!r}exceptasyncio.TimeoutError:outcome抛出了 TimeoutErrorprint(f 期望: 1.00s 抛 TimeoutError)print(f 实测:{time.perf_counter()-t0:.2f}s,{outcome})asyncdefcase2_retry_budget()-None:print(用例 2: 超时放在重试里面, 总耗时失控)attempts0asyncdefcall_once()-str:nonlocalattempts attempts1awaitasyncio.sleep(PER_CALL_TIMEOUT0.2)# 每次都刚好超时returnokasyncdefwith_retry()-str:last:Exception|NoneNonefor_inrange(4):try:returnawaitasyncio.wait_for(call_once(),timeoutPER_CALL_TIMEOUT)exceptasyncio.TimeoutErrorasexc:lastexcraiselast# type: ignore[misc]t0time.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)raisetaskasyncio.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)returndoneasyncdefcase4_executor_not_cancellable()-None:print(用例 4: executor 里的阻塞任务取消不掉, 会连累后面的任务)loopasyncio.get_running_loop()# 只给 1 个线程: 线上连接池、线程池被打满就是这个效果withThreadPoolExecutor(max_workers1)aspool:t0time.perf_counter()fut_aloop.run_in_executor(pool,blocking_io,3.0)try:awaitasyncio.wait_for(fut_a,timeout1.0)exceptasyncio.TimeoutError:print(f 任务A:{time.perf_counter()-t0:.2f}s 时超时返回, f但 worker 线程还在跑 blocking_io)t1time.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)returnlate-resultexceptasyncio.CancelledError:logs.append(downstream-cancelled)raiset0time.perf_counter()try:awaitasyncio.wait_for(asyncio.shield(slow()),timeout0.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:float0.4)-str:外层给总预算, 内层给单次上限, 取消必须能穿透。deadlineasyncio.get_running_loop().time()total_budgetasyncdefonce()-str:remainingdeadline-asyncio.get_running_loop().time()ifremaining0:raiseasyncio.TimeoutError(总预算已耗尽)asyncwithasyncio.timeout(min(PER_CALL_TIMEOUT,remaining)):awaitasyncio.sleep(delay)returnfbody-of-{url}last:Exception|NoneNoneattempts0for_inrange(3):attempts1try:bodyawaitonce()fetch_with_budget.last_attemptsattemptsreturnbodyexceptasyncio.TimeoutErrorasexc:lastexcexceptasyncio.CancelledError:raisefetch_with_budget.last_attemptsattemptsraiselast# type: ignore[misc]fetch_with_budget.last_attempts0asyncdefcase5_correct_pattern()-None:print(用例 5: 带总预算的正确写法)t0time.perf_counter()bodyawaitfetch_with_budget(https://example.com/api,total_budget0.6,delay0.15)print(f 场景A 下游 0.15s / 预算 0.6s: f{time.perf_counter()-t0:.2f}s 成功,{body!r})t0time.perf_counter()try:awaitfetch_with_budget(https://example.com/down,total_budget0.5,delay1.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 retrievedasyncio.get_running_loop().set_exception_handler(lambdaloop,ctx:None)print(fPython{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的结果过一遍凡是超时参数写在重试循环里面的按本文的三层预算改成总预算版本。这一步改动通常不超过二十行但能挡掉大部分日志突然不动的故障。