news 2026/9/24 20:51:14

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 tasks

04: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-cancel

CancelledError从 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就能跑,无第三方依赖:

展开完整代码(约 200 行)
#!/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的结果过一遍,凡是超时参数写在重试循环里面的,按本文的三层预算改成总预算版本。这一步改动通常不超过二十行,但能挡掉大部分"日志突然不动"的故障。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/24 20:50:42

AI工程全景地图:从模型到系统落地,六大能力域与工程实践解析

去年我在一个制造业客户的会议室里&#xff0c;听他们IT负责人讲了一个特别典型的事&#xff1a;算法团队花三个月训练了一个设备故障预测模型&#xff0c;准确率看着不错&#xff0c;但真到了产线上&#xff0c;数据接入要重新写管道&#xff0c;特征口径跟早会报表对不上&…

作者头像 李华
网站建设 2026/9/24 20:50:15

ASP+AJAX在老旧系统中的实战应用与避坑指南

1. 这不是“过时技术”的怀旧表演&#xff0c;而是真实生产环境里仍在呼吸的Web骨架你点开这个标题&#xff0c;心里可能已经浮现出几个问号&#xff1a;ASP&#xff1f;那个用VBScript写<% Response.Write "Hello World" %>的古董&#xff1f;AJAX&#xff1f…

作者头像 李华
网站建设 2026/9/24 20:49:58

AI Agent + Tabular Editor:让大模型直接操作Power BI模型的实战指南

做Power BI模型开发的朋友&#xff0c;对Tabular Editor这个名字应该不陌生。最近半年我把这个工具和AI Agent组合到一起&#xff0c;摸索了一套“让大模型直接动手改Power BI模型”的开发工作流&#xff0c;今天把整套思路和踩坑记录完整聊一遍。无论你是刚开始接触Power BI建…

作者头像 李华
网站建设 2026/9/24 20:49:05

AI生成PPT工具怎么选?Agent路线实现专业级排版与设计

1. 为什么“专业级PPT”这件事&#xff0c;AI工具的选择比努力更重要做PPT这件事&#xff0c;几乎每个职场人都绕不开。不管你是做技术方案汇报、产品路演、年终总结&#xff0c;还是给学生上课、参加创业比赛&#xff0c;PPT都是绕不过去的一道坎。我见过太多人&#xff0c;内…

作者头像 李华
网站建设 2026/9/24 20:48:55

AI室内设计会改结构吗?四款工具实测与避坑指南

1. 从一张户型图说起&#xff1a;AI室内设计到底动了什么很多人第一次用AI做室内设计&#xff0c;心里都揣着同一个疑问&#xff1a;我把户型图丢进去&#xff0c;它会不会自作主张把承重墙砸了、把窗户挪了、把卫生间改到客厅中间&#xff1f;这个担心不是多余的。我前后用四款…

作者头像 李华