如何在Asyncio中根据订阅Task状态控制另一个Task的启停?
无限循环赛事任务流程实现方案
要实现仅在首次进入in progress状态或取消后首次进入该状态时启动task_2,核心是通过状态标记控制任务的启动时机,具体实现思路如下:
- 引入状态标记:在每个赛事周期内初始化两个变量:
task_2(存储subscription2任务实例)和task_2_active(布尔值,标记subscription2是否正在运行)。 - 修改subscription1逻辑:当收到
in progress状态时,仅在task_2_active为False时才创建并启动task_2,同时将task_2_active设为True。 - 重置状态标记:在
task_2的代码中加入finally块,无论任务是正常完成还是被取消,都将task_2_active重置为False,确保下次进入in progress状态时能重新启动任务。
示例代码:
import asyncio async def execute_subscription2(): # subscription2核心逻辑:记录输出等操作 try: while True: # 模拟接收赛事数据 data = await receive_ws_data() print(f"subscription2输出: {data}") # 赛事完成时退出循环 if data.get('status') == 'completed': break finally: # 任务结束/取消后重置状态标记 global task_2_active task_2_active = False async def execute_subscription1(ws): global task_2, task_2_active while True: # 接收subscription1的状态数据 status_data = await ws.receive_json() current_status = status_data.get('status') if current_status == 'in progress': # 仅当任务未运行时启动subscription2 if not task_2_active: task_2 = asyncio.create_task(execute_subscription2()) task_2_active = True print("启动subscription2任务") elif current_status == 'completed': # 赛事完成,取消subscription2(若存在且未结束) if task_2 and not task_2.done(): task_2.cancel() await task_2 # 等待取消完成,避免警告 break # 退出subscription1循环,准备关闭连接 async def weekly_event_loop(): while True: # 每周检查赛事状态 print("每周检查赛事") event_start_time = await check_event_start_time() # 休眠至赛事启动 await asyncio.sleep(event_start_time - asyncio.get_event_loop().time()) print("赛事启动,建立WebSocket连接") # 建立WebSocket连接(实际项目替换为真实WS库,如aiohttp) ws = await connect_websocket() # 初始化周期内的任务状态 global task_2, task_2_active task_2 = None task_2_active = False try: await execute_subscription1(ws) finally: # 关闭连接,进入下一周循环 await ws.close() print("关闭WebSocket连接,等待下一周检查") await asyncio.sleep(604800) # 休眠一周 # 辅助模拟函数 async def check_event_start_time(): # 模拟返回赛事启动时间(当前时间+10秒,用于测试) return asyncio.get_event_loop().time() + 10 async def connect_websocket(): class MockWS: async def receive_json(self): await asyncio.sleep(1) return {'status': 'in progress'} async def close(self): print("WebSocket连接已关闭") return MockWS() async def receive_ws_data(): await asyncio.sleep(2) return {'status': 'completed'} if __name__ == "__main__": asyncio.run(weekly_event_loop())
取消不存在的任务的异常问题
针对不同场景,取消操作的行为和异常情况如下:
- 取消
None类型的任务变量:会触发AttributeError,因为None没有cancel()方法。 - 取消已完成的任务:
cancel()方法会返回False,不会触发任何异常(任务已结束,无法被取消)。 asyncio.CancelledError:该异常是在任务内部触发的——当任务在await阶段被取消时,任务代码中会抛出此异常,而非调用cancel()方法时触发。asyncio.InvalidStateError:此异常通常在任务处于不允许的状态时操作触发(例如尝试获取未完成任务的结果),取消不存在/已完成的任务不会触发该异常。
示例验证代码:
import asyncio async def test_cancel_scenarios(): # 场景1:取消None变量 task = None try: task.cancel() except AttributeError as e: print(f"取消None变量触发异常: {e}") # 输出: 'NoneType' object has no attribute 'cancel' # 场景2:取消已完成的任务 task = asyncio.create_task(asyncio.sleep(0.1)) await asyncio.sleep(0.2) print(f"任务是否完成: {task.done()}") # 输出: True cancel_result = task.cancel() print(f"取消已完成任务的返回值: {cancel_result}") # 输出: False # 场景3:取消运行中的任务(触发内部CancelledError) task = asyncio.create_task(asyncio.sleep(10)) cancel_result = task.cancel() print(f"取消运行中任务的返回值: {cancel_result}") # 输出: True try: await task except asyncio.CancelledError: print("任务内部触发CancelledError") asyncio.run(test_cancel_scenarios())
内容的提问来源于stack exchange,提问作者HJA24
相关产品推荐
相关产品推荐

