如何将控制权交还给事件循环,实现asyncio任务的并发取消?
问题代码
import asyncio async def long_running_task(): try: while True: # Your long-running processing if asyncio.current_task().cancelled(): raise asyncio.CancelledError except asyncio.CancelledError: print("Task was cancelled") async def main(): task = asyncio.create_task(long_running_task()) await asyncio.sleep(5) # Let the task run for 5 seconds task.cancel() # Cancel the task print("cancelled") asyncio.run(main())
问题描述
上述代码中,task.cancel()始终无法生效,因为Python异步采用单线程模型,而long_running_task()里没有执行任何await操作,事件循环被这个任务持续占用,根本轮不到执行取消逻辑。
实际业务场景中,我存在一些阻塞式调用,希望任务被取消后就不再执行这些调用,但现在任务一直处于运行状态,没法被取消。请问该怎么实现正确的并发?
解决方法
针对异步任务中的阻塞操作,核心有两种处理思路:
1. 将阻塞调用转移到线程池/进程池执行
异步事件循环无法处理纯阻塞的同步代码,因此可以用asyncio.to_thread()(Python 3.9+)或loop.run_in_executor()把阻塞操作丢到线程池,让事件循环能继续调度其他任务,包括处理取消信号。
示例代码:
import asyncio import time # 模拟业务中的阻塞调用 def blocking_business_call(): time.sleep(1) # 模拟耗时阻塞操作 return "业务处理完成" async def long_running_task(): try: while True: # 把阻塞操作交给线程池,await让出事件循环控制权 result = await asyncio.to_thread(blocking_business_call) print(result) # await操作会自动响应取消信号,无需手动检查 except asyncio.CancelledError: print("任务已被取消") async def main(): task = asyncio.create_task(long_running_task()) await asyncio.sleep(5) task.cancel() await task # 等待任务处理取消逻辑 print("已触发取消") asyncio.run(main())
2. 在阻塞逻辑间隙主动检查取消状态
如果无法将阻塞操作放到线程池(比如必须在当前线程执行),可以把大的阻塞逻辑拆分成多个小步骤,每完成一步就主动检查任务的取消状态。
示例代码:
import asyncio async def long_running_task(): try: while True: # 把重型阻塞逻辑拆分成小批次执行 for _ in range(1000000): # 模拟小量阻塞计算 pass # 每完成一批就检查取消状态 if asyncio.current_task().cancelled(): raise asyncio.CancelledError print("完成一批次处理") except asyncio.CancelledError: print("任务已被取消") async def main(): task = asyncio.create_task(long_running_task()) await asyncio.sleep(1) task.cancel() await task print("已触发取消") asyncio.run(main())
关键注意点
- 异步任务要能响应取消,必须通过
await让出事件循环控制权,或者主动检查取消状态。 - IO密集型阻塞操作优先用线程池,线程切换开销小;CPU密集型阻塞操作建议用进程池(通过
concurrent.futures.ProcessPoolExecutor),规避GIL限制。
内容的提问来源于stack exchange,提问作者Tom Huntington
相关产品推荐
相关产品推荐

