Python中能否实现异步任务并行计时器?处理协程超时暂停恢复场景
实现asyncio后台任务超时触发当前协程暂停处理后恢复工作
可以实现你的需求,核心是通过并发监听任务状态和超时事件,让当前协程在执行有效工作的同时,能及时响应超时触发,处理后再恢复工作。asyncio.wait_for因为会阻塞等待,所以不适用,我们可以用以下两种方案实现:
方案一:拆分有效工作为可中断步骤(推荐)
将你的有效工作拆分为多个小步骤,结合asyncio.wait的FIRST_COMPLETED模式,同时监听工作步骤完成、超时触发、长任务完成三个事件,确保任何事件发生时都能立即响应。
import asyncio async def long_task(): try: print("Long task started") await asyncio.sleep(10) # 模拟长时间操作 return "Long task completed" except asyncio.CancelledError: print("Long task cancelled") raise async def useful_work_step(step_num): # 模拟可拆分的有效工作步骤 print(f"Executing useful work step {step_num}") await asyncio.sleep(0.5) # 模拟单步工作耗时 return f"Step {step_num} done" async def main(): long_task_handle = asyncio.create_task(long_task()) timeout_seconds = 3 # 创建超时监控任务 timeout_task = asyncio.create_task(asyncio.sleep(timeout_seconds)) timeout_handled = False # 执行多步骤有效工作 for step in range(1, 10): # 等待第一个完成的事件:工作步骤/超时/长任务完成 done, pending = await asyncio.wait( [useful_work_step(step), timeout_task, long_task_handle], return_when=asyncio.FIRST_COMPLETED ) if timeout_task in done: # 处理超时逻辑(仅执行一次) if not timeout_handled: print(f"Timeout reached after {timeout_seconds}s! Processing timeout...") long_task_handle.cancel() try: await long_task_handle except asyncio.CancelledError: pass timeout_handled = True timeout_task = None # 标记超时任务已处理 elif long_task_handle in done: # 长任务提前完成,清理超时任务并退出循环 result = await long_task_handle print(f"Long task finished early: {result}") if not timeout_task.done(): timeout_task.cancel() break else: # 当前工作步骤完成,继续下一个步骤 result = done.pop().result() print(result) # 清理剩余未完成的任务 if not long_task_handle.done(): long_task_handle.cancel() await long_task_handle if timeout_task and not timeout_task.done(): timeout_task.cancel() await timeout_task asyncio.run(main())
方案一关键点
- 拆分工作为小步骤,让事件循环有机会处理超时和长任务事件
asyncio.wait的FIRST_COMPLETED模式确保优先响应最先发生的事件- 通过
timeout_handled避免重复处理超时逻辑 - 超时触发后可按需取消长任务,处理完后继续执行剩余工作
方案二:持续工作中定期让出控制权
如果你的有效工作无法拆分为步骤,可以在持续工作中定期调用await asyncio.sleep(0)让出事件循环控制权,配合asyncio.Event监听超时触发信号。
import asyncio async def long_task(): try: print("Long task started") await asyncio.sleep(10) return "Long task completed" except asyncio.CancelledError: print("Long task cancelled") raise async def main(): long_task_handle = asyncio.create_task(long_task()) timeout_seconds = 3 timeout_event = asyncio.Event() # 超时监控协程 async def timeout_monitor(): await asyncio.sleep(timeout_seconds) timeout_event.set() asyncio.create_task(timeout_monitor()) timeout_handled = False # 模拟持续有效工作 print("Starting continuous useful work...") work_count = 0 while True: # 模拟工作操作(这里用计数代替实际工作) work_count += 1 # 每处理一定量的工作后,让出控制权检查事件 if work_count % 100000 == 0: print(f"Work in progress: processed {work_count} units") await asyncio.sleep(0) # 检查超时事件 if timeout_event.is_set() and not timeout_handled: print(f"Timeout reached after {timeout_seconds}s! Processing timeout...") long_task_handle.cancel() try: await long_task_handle except asyncio.CancelledError: pass timeout_handled = True # 模拟工作结束条件 if work_count >= 10000000: break print("Continuous useful work finished") # 清理剩余任务 if not long_task_handle.done(): long_task_handle.cancel() await long_task_handle asyncio.run(main())
方案二关键点
- 用
asyncio.Event作为超时触发的信号 await asyncio.sleep(0)不会实际休眠,但会让出事件循环,让超时监控协程有机会执行- 持续工作中定期检查事件,触发后处理超时逻辑,之后恢复工作
内容的提问来源于stack exchange,提问作者Garry
相关产品推荐
相关产品推荐

