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

Python异步编程取消任务遇RecursionDepthError求解决方案

Python异步任务取消递归错误与取消失败问题解决

问题根源分析

  1. 递归错误触发原因:cleanup()协程通过run_coroutine_threadsafe提交到事件循环后,会成为当前循环中的活跃任务。执行asyncio.all_tasks(self.loop)时,这个正在运行的cleanup任务会被纳入待取消列表。尝试取消自身并通过await asyncio.gather等待完成时,会触发任务取消逻辑的递归调用,最终超出Python递归深度限制。
  2. 任务取消失败原因:task.cancel()仅发送取消请求,任务不会立即进入已取消状态——只有当任务下次被调度执行到await等挂起点时,才会响应取消并抛出CancelledError。直接在cancel()后断言task.cancelled()为True,忽略了异步任务状态变更的延迟性;同时取消当前正在运行的任务本身,也无法让状态立刻变更。

修复方案

修改代码核心逻辑,排除当前任务并调整取消后的等待逻辑:

import asyncio
import threading
import logging

console_logger = logging.getLogger("console_logger")
console_logger.setLevel(logging.DEBUG)
handler = logging.StreamHandler()
console_logger.addHandler(handler)

class Example:
    def __init__(self):
        self.loop = asyncio.new_event_loop()
        self.thread = threading.Thread(target=self.run_loop, args=(self.loop,))
        self.thread.start()
        self.websocket = None  # 模拟websocket连接对象

    def run_loop(self, loop):
        asyncio.set_event_loop(loop)
        loop.run_forever()  # 改为持续运行事件循环,符合实际项目场景

    def close(self):
        """关闭WebSocket连接并清理所有异步任务"""

        def stop_loop(loop: asyncio.AbstractEventLoop):
            loop.stop()
            loop.close()

        async def cleanup():
            # 处理WebSocket关闭逻辑
            if self.websocket and not self.websocket.closed:
                await self.websocket.close()
            
            # 获取当前执行的cleanup任务,排除自身避免递归取消
            current_task = asyncio.current_task(self.loop)
            tasks = [
                t for t in asyncio.all_tasks(self.loop) 
                if not t.done() and t is not current_task
            ]
            
            console_logger.debug(f"找到 {len(tasks)} 个待取消任务.")

            for task in tasks:
                console_logger.debug(f"发起取消请求: {task}")
                task.cancel()

            console_logger.debug("等待任务清理完成")
            # 等待所有任务处理取消,return_exceptions=True避免异常中断流程
            await asyncio.gather(*tasks, return_exceptions=True)
            console_logger.debug("任务清理完成")
            self.loop.call_soon_threadsafe(stop_loop, self.loop)

        future = asyncio.run_coroutine_threadsafe(cleanup(), self.loop)
        try:
            future.result()
        except Exception as e:
            console_logger.debug(f"清理过程出现异常: {e}")
        self.thread.join()


# 示例运行
example = Example()
example.close()

关键修改说明

  • 排除当前任务:通过asyncio.current_task(self.loop)获取正在执行的cleanup任务,从待取消列表中剔除,彻底避免递归取消的问题。
  • 调整事件循环启动方式:将run_loop中的loop.run_until_complete(self.close())改为loop.run_forever(),符合实际项目中事件循环长期运行的逻辑,避免初始化阶段就触发关闭流程。
  • 合理处理取消等待:await asyncio.gather(*tasks, return_exceptions=True)会等待所有任务响应取消请求,即使任务抛出CancelledError也不会中断清理流程。
  • 取消状态断言移除:取消请求发送后无需立即断言状态,任务会在合适的时机响应取消;若需验证取消结果,可在gather完成后遍历任务检查task.cancelled()。

额外注意事项

  • 任务取消的局限性:如果被取消的任务是无挂起点的无限循环,cancel()无法终止它。这种情况下需要在任务逻辑中定期检查asyncio.current_task().cancelled(),主动退出循环。
  • 线程安全保障:跨线程操作事件循环时,必须使用run_coroutine_threadsafe和call_soon_threadsafe,原代码的这部分实现是正确的。

内容的提问来源于stack exchange,提问作者Sayan Bose

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:44:57