使用asyncio.wait_for+shield仍触发超时的问题排查
问题描述
我用Python结合aioboto3写了异步下载S3文件的代码,运行时部分文件能成功下载解压,但突然触发asyncio.TimeoutError导致任务取消,而且耗时还不到一分钟——我试过调大超时值,完全没用。想搞懂问题出在哪,为什么asyncio.shield没能阻止任务超时?
代码实现
import boto3 import botocore.exceptions as bexc import asyncio import aioboto3 from aiobotocore.session import get_session import gzip import utils.processor as processor async def download_s3_async(self, file): session = aioboto3.Session(profile_name=self.current_profile) async with session.client("s3") as client: start_log = processor.cf_get_start_time_from_filename(file.split("/")[3]) end_log = processor.cf_get_end_time_from_filename(file.split("/")[3]) s3_obj = await client.get_object(Bucket=self.cf_bucket, Key=file) task = s3_obj["Body"].read() try: cat = await asyncio.wait_for(asyncio.shield(task), timeout=60000000) except asyncio.TimeoutError: print('Gave up waiting, task canceled') res = gzip.decompress(cat).decode("utf-8") print("File with {}-{} timestamp was downloaded and decompressed".format(start_log, end_log)) return res async def iterate(self, files): result = await asyncio.gather(*map(self.download_s3_async, files)) return result def call_async_wrapper(self, s_time, e_time): files = self.get_s3_objects_list(s_time, e_time) loop = asyncio.get_event_loop() results = loop.run_until_complete(self.iterate(files)) return results def get_waf_actors_info(self, s_date, s_time, e_time): self.current_dist = "somepath/path/{}/".format(str(s_date).replace("-", "")) self.call_async_wrapper(s_date, s_time, e_time)
运行日志
... File with 00:01:29-00:01:59 timestamp was downlaoded and decompressed File with 00:00:24-00:00:54 timestamp was downlaoded and decompressed File with 00:03:33-00:04:03 timestamp was downlaoded and decompressed File with 00:07:01-00:07:31 timestamp was downlaoded and decompressed File with 00:01:01-00:01:31 timestamp was downlaoded and decompressed File with 00:06:14-00:06:44 timestamp was downlaoded and decompressed File with 00:03:59-00:04:29 timestamp was downlaoded and decompressed File with 00:03:02-00:03:32 timestamp was downlaoded and decompressed File with 00:09:06-00:09:36 timestamp was downlaoded and decompressed File with 00:04:12-00:04:42 timestamp was downlaoded and decompressed File with 00:02:32-00:03:02 timestamp was downlaoded and decompressed File with 00:02:23-00:02:53 timestamp was downlaoded and decompressed File with 00:00:08-00:00:38 timestamp was downlaoded and decompressed File with 00:09:31-00:10:01 timestamp was downlaoded and decompressed File with 00:00:23-00:00:53 timestamp was downlaoded and decompressed Gave up waiting, task canceled
问题分析与解决
1. 你误解了asyncio.shield的作用
asyncio.shield的功能是防止被外部任务取消,但它不会改变asyncio.wait_for本身的超时逻辑——wait_for到了设定时间依然会抛出TimeoutError,只是不会把取消信号传递给被shield的任务。但你这里设置的timeout=60000000(约16666小时),显然不是这个超时触发的,所以问题不在这。
2. 真正的元凶:aiobotocore内置的HTTP超时
aioboto3的底层依赖aiobotocore,它的客户端默认有自己的HTTP超时设置(connect_timeout默认60秒,read_timeout默认60秒)。当下载大文件或者网络较慢时,底层HTTP客户端会先触发超时,这个异常会让s3_obj["Body"].read()抛出错误,而这个错误被wait_for捕获后,表现出来像是asyncio.TimeoutError,和你设置的wait_for超时完全无关。
3. 代码里的额外问题:任务创建方式错误
你写的task = s3_obj["Body"].read()只是把一个协程对象赋值给了task,并没有真正创建异步任务。asyncio.shield需要包裹的是asyncio.Task对象,否则起不到任何保护作用。正确的写法应该是:
task = asyncio.create_task(s3_obj["Body"].read())
4. 解决步骤
- 调大aiobotocore的超时配置:创建S3客户端时,显式设置超时参数,覆盖默认值:
import botocore.config async with session.client( "s3", config=botocore.config.Config( read_timeout=300, # 5分钟,根据文件大小调整 connect_timeout=60 ) ) as client: - 修正任务创建方式:用
asyncio.create_task创建任务后再用shield包裹。 - 区分异常类型:捕获botocore的超时异常,避免和asyncio的TimeoutError混淆:
try: cat = await asyncio.wait_for(asyncio.shield(task), timeout=60000000) except asyncio.TimeoutError: print('Asyncio等待超时') except bexc.ReadTimeoutError: print('S3下载读取超时') - 复用session:不要每次下载都创建新的
aioboto3.Session,可以把session作为类属性复用,提升性能。
内容的提问来源于stack exchange,提问作者ShHolmes
相关产品推荐
相关产品推荐

