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

Python自修复任务组实现:如何自动重启失败任务?

优化Python异步任务的自修复任务管理器实现

我想在Python中实现一个能够重启失败任务的任务管理器,以下是我构思的实现代码,但感觉像是临时的权宜之计。请问有没有更优的方式实现这种“自修复”任务组模式?

用户提供的初始实现代码:

import asyncio, random

async def noreturn(_arg):
    while True:
        await asyncio.sleep(1)
        
        if random.randint(0, 10) % 10 == 0:
            raise random.choice((RuntimeError, ValueError, TimeoutError))


async def main():
    taskmap: dict[int, asyncio.Task] = {}

    for i in range(10):
        taskmap[i] = asyncio.create_task(noreturn(i))
    
    while True:
        for arg, task in taskmap.items():
            if task.done():
                # Task died
                taskmap[arg] = asyncio.create_task(noreturn(arg))

        await asyncio.sleep(1)


if __name__ == "__main__":
    asyncio.run(main())

初始代码的核心问题

  • 轮询检查任务状态存在延迟:任务失败后最多要等1秒才会重启,无法即时响应
  • 缺少异常记录:任务失败的原因无法追踪,不利于调试
  • 逻辑耦合度高:任务重启逻辑和主循环混在一起,扩展和维护成本高

优化方案一:任务包装器+回调实现即时重启

利用asyncio.Task.add_done_callback实现任务失败后的即时重启,同时解耦重启逻辑,增加异常日志记录。

实现代码

import asyncio
import random
import logging

# 配置日志,方便追踪任务状态
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

async def noreturn(_arg):
    while True:
        await asyncio.sleep(1)
        if random.randint(0, 10) % 10 == 0:
            raise random.choice((RuntimeError, ValueError, TimeoutError))

def get_restart_callback(task_id, task_func, task_map):
    """生成任务失败后的重启回调函数"""
    def callback(future):
        # 获取并记录任务异常信息
        exc = future.exception()
        if exc:
            logger.warning(f"任务 {task_id} 失败,异常类型: {type(exc).__name__},详情: {exc}")
        
        # 重启任务并更新任务映射
        new_task = asyncio.create_task(task_func(task_id))
        new_task.add_done_callback(callback)
        task_map[task_id] = new_task
        logger.info(f"任务 {task_id} 已重启")
    return callback

async def main():
    task_map = {}
    task_count = 10

    # 初始化所有任务并绑定回调
    for i in range(task_count):
        task = asyncio.create_task(noreturn(i))
        task.add_done_callback(get_restart_callback(i, noreturn, task_map))
        task_map[i] = task

    # 保持主任务运行,避免程序退出
    await asyncio.Event().wait()

if __name__ == "__main__":
    asyncio.run(main())

优化点说明

  • 即时响应:任务失败后立即触发重启,无需轮询等待
  • 异常追踪:记录任务失败的异常类型和详情,便于问题排查
  • 逻辑解耦:重启逻辑封装在独立函数中,与任务业务逻辑分离

优化方案二:封装成TaskManager类实现精细化管理

如果需要更复杂的管理逻辑(比如动态增删任务、限制重启次数、设置重启延迟),可以封装成类,提升代码的可维护性和扩展性。

实现代码

import asyncio
import random
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

async def noreturn(_arg):
    while True:
        await asyncio.sleep(1)
        if random.randint(0, 10) % 10 == 0:
            raise random.choice((RuntimeError, ValueError, TimeoutError))

class TaskManager:
    def __init__(self):
        self._tasks = {}  # 存储任务ID与对应Task对象的映射

    def _build_restart_callback(self, task_id, task_func, max_restarts=None, restart_delay=0):
        """生成带重启限制和延迟的回调函数"""
        restart_count = 0

        def callback(future):
            nonlocal restart_count
            exc = future.exception()
            if exc:
                logger.warning(f"任务 {task_id} 失败,异常: {type(exc).__name__}: {exc}")

                # 检查是否达到最大重启次数
                if max_restarts is not None and restart_count >= max_restarts:
                    logger.error(f"任务 {task_id} 已达最大重启次数 {max_restarts},停止重启")
                    del self._tasks[task_id]
                    return

                restart_count += 1
                logger.info(f"任务 {task_id} 准备重启,当前重启次数: {restart_count}")

                # 带延迟重启任务
                async def delayed_restart():
                    await asyncio.sleep(restart_delay)
                    new_task = asyncio.create_task(task_func(task_id))
                    new_task.add_done_callback(callback)
                    self._tasks[task_id] = new_task
                asyncio.create_task(delayed_restart())
            else:
                logger.info(f"任务 {task_id} 正常完成")
                del self._tasks[task_id]
        return callback

    def add_task(self, task_id, task_func, max_restarts=None, restart_delay=0):
        """添加需要管理的任务"""
        if task_id in self._tasks:
            logger.warning(f"任务 {task_id} 已存在,跳过添加")
            return
        task = asyncio.create_task(task_func(task_id))
        task.add_done_callback(self._build_restart_callback(task_id, task_func, max_restarts, restart_delay))
        self._tasks[task_id] = task
        logger.info(f"任务 {task_id} 已添加")

    def remove_task(self, task_id):
        """移除指定任务"""
        if task_id in self._tasks:
            self._tasks[task_id].cancel()
            del self._tasks[task_id]
            logger.info(f"任务 {task_id} 已移除")

    async def run_forever(self):
        """保持管理器持续运行"""
        await asyncio.Event().wait()

async def main():
    manager = TaskManager()
    # 添加10个任务,每个任务最多重启5次,重启前延迟0.5秒
    for i in range(10):
        manager.add_task(i, noreturn, max_restarts=5, restart_delay=0.5)
    await manager.run_forever()

if __name__ == "__main__":
    asyncio.run(main())

进阶优化点

  • 动态任务管理:支持随时添加/删除任务
  • 重启控制:可设置最大重启次数,避免无限重启故障任务
  • 延迟重启:防止任务频繁失败导致资源耗尽
  • 状态统一管理:通过类属性集中管理所有任务的运行状态

内容的提问来源于stack exchange,提问作者Joe Boyle

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 10:52:43