Python Asyncio任务超时与重试设计实现请求
Python Asyncio 多任务独立超时重试+全局总超时实现
需求说明
- 每个异步任务配置独立的超时时间与重试次数,超时后自动重试,直至成功或用尽重试次数
- 所有任务(含重试流程)必须在全局总超时时间内完成,超时未完成的任务直接返回
None - 任务并发执行,最终返回三元组结果:成功任务返回对应字符串,失败/被中断任务返回
None
完整示例代码
import asyncio import typing as tp import random async def my_func(wait_time: int) -> str: random_number = random.random() random_time = wait_time - random_number if random.random() < 0.5 else wait_time + random_number print(f"[{asyncio.current_task().get_name()}] 预计等待 {wait_time}{random_time:+} 秒") await asyncio.sleep(wait_time) return f"[{asyncio.current_task().get_name()}] 实际等待 {wait_time}{random_time:+} 秒" async def run_with_retry_and_timeout( func: tp.Callable[..., tp.Awaitable[str]], func_args: tp.Tuple, timeout: float, max_retry: int, task_name: str ) -> tp.Optional[str]: """封装单个任务的重试与独立超时逻辑""" asyncio.current_task().set_name(task_name) for attempt in range(max_retry + 1): # 包含首次执行,共max_retry+1次执行机会 try: # 单次执行的超时检查 return await asyncio.wait_for(func(*func_args), timeout=timeout) except asyncio.TimeoutError: if attempt >= max_retry: print(f"[{task_name}] 重试{max_retry}次后仍超时,放弃执行") return None print(f"[{task_name}] 第{attempt+1}次执行超时,进行重试") except asyncio.CancelledError: # 被全局超时中断 print(f"[{task_name}] 被全局超时终止") return None except Exception as e: # 处理任务执行中的其他异常(如业务错误) print(f"[{task_name}] 第{attempt+1}次执行出错: {str(e)},进行重试") if attempt >= max_retry: print(f"[{task_name}] 重试{max_retry}次后仍出错,放弃执行") return None async def main() -> tp.Tuple[tp.Optional[str], tp.Optional[str], tp.Optional[str]]: # 任务配置:(函数, 参数, 独立超时, 重试次数, 任务名) task_configs = [ (my_func, (1,), 1.2, 4, "task1"), (my_func, (2,), 2.2, 3, "task2"), (my_func, (3,), 3.2, 2, "task3"), ] total_timeout = 5 # 创建带重试和独立超时的任务列表 tasks = [ run_with_retry_and_timeout(func, args, timeout, retry, name) for func, args, timeout, retry, name in task_configs ] # 全局超时控制 gather_task = asyncio.gather(*tasks, return_exceptions=True) try: results = await asyncio.wait_for(gather_task, timeout=total_timeout) except asyncio.TimeoutError: print("\n全局总超时触发,终止所有未完成任务") gather_task.cancel() results = await gather_task # 处理结果:将异常/取消状态转为None,保留成功结果 processed_results = [res if isinstance(res, str) else None for res in results] return tuple(processed_results) if __name__ == "__main__": task1_res, task2_res, task3_res = asyncio.run(main()) print("\n最终结果:") print(f"task1: {task1_res}") print(f"task2: {task2_res}") print(f"task3: {task3_res}")
关键逻辑说明
- 单任务封装函数:
run_with_retry_and_timeout负责处理单个任务的重试流程与独立超时,同时捕获全局超时触发的CancelledError,返回None - 全局超时控制:用
asyncio.wait_for包裹asyncio.gather,确保所有任务的总执行时间不超过设定值;全局超时触发时,取消所有未完成任务并等待收尾 - 结果处理:通过
return_exceptions=True让gather捕获所有任务的异常,最终将异常/取消状态转为None,保留成功任务的返回值
内容的提问来源于stack exchange,提问作者user3225309
相关产品推荐
相关产品推荐

