如何取消长时间运行的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的线程中执行的:
- Python线程没有安全的强制终止机制,一旦线程开始运行同步函数,
asyncio无法直接中断它; concurrent.futures.Future.cancel()对已经启动的线程任务也无效,只能取消未开始的任务。
解决方案:改用ProcessPoolExecutor替代线程池
进程支持强制终止,因此用进程池执行同步任务时,超时后可以直接终止对应进程,从而中断_work的执行。
修改步骤
- 替换线程池为进程池
- 跟踪每个异步任务对应的进程池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
注意事项
- 进程池的开销比线程池大,因为进程间通信存在额外成本,需根据实际场景评估性能;
- 进程被强制终止可能导致外部资源(如数据库连接)泄漏,需确认外部驱动能处理这种异常终止情况;
- 如果任务依赖共享内存对象(如线程安全的缓存),进程池可能无法直接使用,需调整实现逻辑。
内容的提问来源于stack exchange,提问作者sanyassh
相关产品推荐
相关产品推荐

