如何正确处理asyncio.TimeoutError?任务异常未捕获问题排查
问题
我用asyncio和aiohttp写了个脚本,要发大量并发HTTP HEAD请求并收集结果。根据asyncio文档,某个HEAD请求触发异常时,会立刻传播到等待gather()的任务。捕获异常时我会取消所有其他任务然后从头重试,这对aiohttp的异常有效,但遇到asyncio.TimeoutError时行为异常。
示例代码:
import aiohttp import asyncio import logging import time async def fetch_content_length(session, url): async with session.head(url) as resp: return resp.content_length async def content_lengths(session, url): endpoint = url + '{}' tasks = [asyncio.create_task(fetch_content_length(session, endpoint.format(i))) for i in range(500)] try: results = [await coro for coro in asyncio.gather(*tasks)] except (aiohttp.ClientResponseError, aiohttp.ClientConnectionError, asyncio.TimeoutError): for t in tasks: t.cancel() raise return results async def all_content_lengths(session): urls = ['https://www.example.com/', 'https://www.example2.com/', 'https://www.example3.com/'] tasks = [asyncio.create_task(content_lengths(session, url)) for url in urls] try: results = [await coro for coro in asyncio.as_completed(tasks)] except (aiohttp.ClientResponseError, aiohttp.ClientConnectionError, asyncio.TimeoutError): for t in tasks: t.cancel() raise return results async def run(): connector = aiohttp.TCPConnector(force_close=True) session = aiohttp.ClientSession(connector=connector, raise_for_status=True) try: for retry in range(1, 6): retry_time = 2 ** retry try: results = await all_content_lengths(session) except aiohttp.ClientResponseError as e: logging.info(f"Failed with {e.status} on {e.request_info.url}") time.sleep(retry_time) except aiohttp.ClientConnectionError as e: logging.info(f"Failed with: {e}") time.sleep(retry_time) except asyncio.TimeoutError: logging.info("Failed with timeout") time.sleep(retry_time) else: return results else: logging.warning("All retries failed") finally: session.close() def main(): results = asyncio.run(run())
出现asyncio.TimeoutError时,程序能重试,但日志里会出现错误:
WARNING: Failed with timeout ERROR: Task exception was never retrieved future: <Task finished name='Task-891068' coro=<content_lengths() done, defined at example.py:10> exception=TimeoutError()> Traceback (most recent call last): File "example.py", line 18, in content_lengths results = await asyncio.gather(*tasks) File "example.py", line 7, in fetch_content_length async with session.head(url) as resp: File "example-project\venv\lib\site-packages\aiohttp\client.py", line 1138, in __aenter__ self._resp = await self._coro File "example-project\venv\lib\site-packages\aiohttp\client.py", line 634, in _request break File "example-project\venv\lib\site-packages\aiohttp\helpers.py", line 721, in __exit__ raise asyncio.TimeoutError from None asyncio.exceptions.TimeoutError
有时重试成功前会连续出现两次这个错误。我搞不懂Task exception was never retrieved的原因,本来以为第一个异常会触发所有其他任务的取消,是不是和其他任务也触发超时异常有关?
解答
问题根源
超时异常触发时机差:第一个请求超时抛出异常后,你在
except块里取消剩余任务,但此时可能已有其他任务也触发了超时——这些任务的异常没被gather()捕获(因为gather()默认在第一个异常抛出后就停止等待),任务结束时异常没人处理,就会触发Task exception was never retrieved报错。gather()的默认行为:默认情况下,gather()遇到第一个异常就会停止等待并传播该异常,剩下的任务会继续运行,但它们的异常不会被收集。哪怕你事后取消任务,已经抛出的异常还是会留在任务对象里,最终触发日志警告。as_completed的隐性问题:在all_content_lengths里用as_completed,当某个content_lengths任务抛出异常后,你取消其他任务,但子任务可能已经抛出异常,没被上层逻辑处理。
解决方法
方法1:让gather()收集所有异常
给gather()加上return_exceptions=True参数,它会把所有任务的结果(包括异常)收集到列表里,不会立刻传播异常。之后你可以遍历结果检查异常,再决定是否取消任务并重试:
async def content_lengths(session, url): endpoint = url + '{}' tasks = [asyncio.create_task(fetch_content_length(session, endpoint.format(i))) for i in range(500)] # 收集所有结果(含异常) results = await asyncio.gather(*tasks, return_exceptions=True) # 检查是否有需要处理的异常 for res in results: if isinstance(res, (aiohttp.ClientResponseError, aiohttp.ClientConnectionError, asyncio.TimeoutError)): # 取消未完成的任务 for t in tasks: if not t.done(): t.cancel() # 抛出第一个异常触发重试 raise res return results
方法2:取消任务后主动处理异常
在except块里取消任务后,主动await每个任务,检索它们的异常,避免出现“未被检索”的情况:
except (aiohttp.ClientResponseError, aiohttp.ClientConnectionError, asyncio.TimeoutError): for t in tasks: t.cancel() # 处理任务可能已抛出的异常 try: await t except (asyncio.CancelledError, aiohttp.ClientResponseError, aiohttp.ClientConnectionError, asyncio.TimeoutError): pass raise
方法3:给任务加异常回调
给每个任务绑定一个回调函数,专门处理未被上层逻辑捕获的异常:
def handle_task_exception(task): try: task.result() except (asyncio.CancelledError, aiohttp.ClientResponseError, aiohttp.ClientConnectionError, asyncio.TimeoutError): pass except Exception as e: logging.error(f"Unexpected task error: {e}") # 创建任务时绑定回调 tasks = [] for i in range(500): task = asyncio.create_task(fetch_content_length(session, endpoint.format(i))) task.add_done_callback(handle_task_exception) tasks.append(task)
额外优化
- 把
time.sleep(retry_time)改成await asyncio.sleep(retry_time),避免阻塞事件循环,符合asyncio的最佳实践。 - 如果不需要按完成顺序处理结果,在
all_content_lengths里用asyncio.gather(*tasks)代替as_completed,异常处理更可控。
内容的提问来源于stack exchange,提问作者embeage

