You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确处理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的原因,本来以为第一个异常会触发所有其他任务的取消,是不是和其他任务也触发超时异常有关?


解答

问题根源

  1. 超时异常触发时机差:第一个请求超时抛出异常后,你在except块里取消剩余任务,但此时可能已有其他任务也触发了超时——这些任务的异常没被gather()捕获(因为gather()默认在第一个异常抛出后就停止等待),任务结束时异常没人处理,就会触发Task exception was never retrieved报错。

  2. gather()的默认行为:默认情况下,gather()遇到第一个异常就会停止等待并传播该异常,剩下的任务会继续运行,但它们的异常不会被收集。哪怕你事后取消任务,已经抛出的异常还是会留在任务对象里,最终触发日志警告。

  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 17:03:20