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

如何取消长时间运行的asyncio.Task?同步外部任务超时终止方案

问题描述

我有一段异步代码,尝试在超时后取消长时间运行的同步任务,但task.cancel()完全无效。运行代码后,即使超时触发,同步任务的_work finished依然会打印出来。我无法修改外部库中的_work函数(它是对接外部数据库的同步驱动),请问该如何正确实现超时后立即终止这些同步操作?

示例代码

import asyncio
import concurrent.futures
import functools
import time


async def run_till_first_success(tasks, timeout=None):
    results = []
    exceptions = []
    while tasks:
        try:
            async with asyncio.timeout(timeout):
                done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED, timeout=timeout)
        except TimeoutError:
            print('timeout is catched')
            for task in tasks:
                task.cancel()  # 这里不起作用!
            return

        # >>> 以下逻辑不影响核心问题
        for task in done:
            if task.exception():
                exceptions.append(task.exception())
            else:
                results.append(task.result())
        tasks = pending
        if results or len(pending) == 0:
            for task in pending:
                task.cancel()
            break
    if not results:
        raise exceptions[0]
    return results[0]
# <<<


class Worker:
    def __init__(self, execute):
        self.execute = execute

    async def work(self):
        return await self.execute(self._work)

    def _work(self):  # 模拟不可修改的外部同步长任务
        print('_work started')
        time.sleep(5)
        print('_work finished')


class LibraryInterfaceClass:
    def __init__(self):
        self._executor = concurrent.futures.ThreadPoolExecutor(10)
        self._workers = [Worker(self._execute) for _ in range(2)]

    async def use(self, timeout=1):
        tasks = {asyncio.create_task(worker.work()) for worker in self._workers}
        return await run_till_first_success(tasks, timeout)

    async def _execute(self, func, *args, **kwargs):
        return await asyncio.get_event_loop().run_in_executor(self._executor, functools.partial(func, *args, **kwargs))


interface = LibraryInterfaceClass()
asyncio.run(interface.use())
print('run is done')

当前输出

_work started
_work started
timeout is catched
run is done
_work finished
_work finished

问题原因

asyncio.Task.cancel()只能中断可被取消的异步等待操作(比如await协程时抛出CancelledError),但你的同步任务_work是在ThreadPoolExecutor的线程中执行的:

  1. Python线程没有安全的强制终止机制,一旦线程开始运行同步函数,asyncio无法直接中断它;
  2. concurrent.futures.Future.cancel()对已经启动的线程任务也无效,只能取消未开始的任务。

解决方案:改用ProcessPoolExecutor替代线程池

进程支持强制终止,因此用进程池执行同步任务时,超时后可以直接终止对应进程,从而中断_work的执行。

修改步骤

  1. 替换线程池为进程池
  2. 跟踪每个异步任务对应的进程池Future,超时后同时取消异步任务和进程池任务

修改后的完整代码:

import asyncio
import concurrent.futures
import functools
import time


async def run_till_first_success(task_future_map, timeout=None):
    results = []
    exceptions = []
    tasks = set(task_future_map.keys())
    while tasks:
        try:
            async with asyncio.timeout(timeout):
                done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED, timeout=timeout)
        except TimeoutError:
            print('timeout is catched')
            for task in tasks:
                task.cancel()
                # 取消对应的进程池Future,强制终止进程
                future = task_future_map[task]
                future.cancel()
            return

        for task in done:
            if task.exception():
                exceptions.append(task.exception())
            else:
                results.append(task.result())
        # 更新待处理的任务和映射关系
        tasks = pending
        task_future_map = {t: f for t, f in task_future_map.items() if t in pending}
        if results or len(pending) == 0:
            for task in pending:
                task.cancel()
                future = task_future_map[task]
                future.cancel()
            break
    if not results:
        raise exceptions[0]
    return results[0]


class Worker:
    def __init__(self, execute):
        self.execute = execute

    async def work(self):
        return await self.execute(self._work)

    def _work(self):
        print('_work started')
        time.sleep(5)
        print('_work finished')


class LibraryInterfaceClass:
    def __init__(self):
        # 替换为ProcessPoolExecutor
        self._executor = concurrent.futures.ProcessPoolExecutor(10)
        self._workers = [Worker(self._execute) for _ in range(2)]

    async def use(self, timeout=1):
        task_future_map = {}
        for worker in self._workers:
            # 先获取进程池的Future
            future = self._execute(worker._work)
            # 创建异步任务并建立映射
            task = asyncio.create_task(future)
            task_future_map[task] = future
        return await run_till_first_success(task_future_map, timeout)

    async def _execute(self, func, *args, **kwargs):
        return await asyncio.get_event_loop().run_in_executor(self._executor, functools.partial(func, *args, **kwargs))


interface = LibraryInterfaceClass()
asyncio.run(interface.use())
print('run is done')

修改后输出

_work started
_work started
timeout is catched
run is done

注意事项

  1. 进程池的开销比线程池大,因为进程间通信存在额外成本,需根据实际场景评估性能;
  2. 进程被强制终止可能导致外部资源(如数据库连接)泄漏,需确认外部驱动能处理这种异常终止情况;
  3. 如果任务依赖共享内存对象(如线程安全的缓存),进程池可能无法直接使用,需调整实现逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 17:31:01