如何在transitions AsyncMachine中执行循环任务且避免死锁?
异步状态机结合MQTT的阻塞问题解决方案
使用transitions库的AsyncMachine结合aiomqtt搭建MQTT控制的异步状态机,基础功能正常,但在B状态添加循环任务(while循环或自反转移)后,状态机卡在B状态,无法响应MQTT的状态变更请求。
问题根源
AsyncMachine的on_enter_*回调会被状态机同步await直到执行完成。如果在on_enter_B里放入无限循环或递归的自反转移,会持续占用事件循环,导致负责接收MQTT消息的receiveMQTT任务无法获得执行时间,进而无法处理状态变更请求。
解决方案
将B状态的循环任务拆分为独立的异步任务,在进入B状态时启动,离开B状态时取消,让事件循环可以同时处理MQTT消息接收和状态循环任务。
修改后的状态机脚本
import asyncio import aiomqtt from transitions.extensions import AsyncMachine import sys import os import logging logging.basicConfig(level=logging.DEBUG) if sys.platform.lower() == "win32" or os.name.lower() == "nt": from asyncio import set_event_loop_policy, WindowsSelectorEventLoopPolicy set_event_loop_policy(WindowsSelectorEventLoopPolicy()) class MQTTStateMachine: states = ['init','A','B',{'name': 'stopped', 'final':True}] def __init__(self,client): self.client = client self.b_task = None # 管理B状态的循环任务 self.machine = AsyncMachine(model=self, states=MQTTStateMachine.states, initial='init') # 补充缺失的状态转移 self.machine.add_transition(trigger='init', source='init', dest='A') self.machine.add_transition(trigger='to_A', source=['B'], dest='A') self.machine.add_transition(trigger='to_B', source=['A'], dest='B') self.machine.add_transition(trigger='stop', source=['A','B'], dest='stopped') async def update_state(self): await self.client.publish("MQTTstatemachine/machine/state", str(self.state)) async def receiveMQTT(self): await self.client.subscribe("MQTTstatemachine/controller/transition") async for message in self.client.messages: if message.topic.matches("MQTTstatemachine/controller/transition"): await self.trigger(message.payload.decode()) async def b_loop_task(self): # B状态的循环任务,独立运行 while self.is_B(): await self.update_state() print("I'm now in state B.") await asyncio.sleep(1) # 让出事件循环,让其他任务执行 async def on_enter_A(self): await self.update_state() print("I'm now in state A.") async def on_enter_B(self): await self.update_state() print("I'm now in state B.") # 启动B状态的循环任务,不await,让事件循环并行处理 self.b_task = asyncio.create_task(self.b_loop_task()) async def on_exit_B(self): # 离开B状态时取消循环任务,避免资源泄漏 if self.b_task and not self.b_task.done(): self.b_task.cancel() try: await self.b_task except asyncio.CancelledError: print("B state loop task cancelled.") async def on_enter_stopped(self): await self.update_state() print("I'm now in state stopped.") async def main(): async with aiomqtt.Client("test.mosquitto.org") as client: MQTTmachine = MQTTStateMachine(client) await MQTTmachine.init() # 启动MQTT接收任务,用gather让任务并行执行 await asyncio.gather( MQTTmachine.receiveMQTT(), ) if __name__ == "__main__": asyncio.run(main())
关键修改点
- 补充状态转移:添加
to_A和to_B的转移规则,对应控制器发送的指令。 - 独立循环任务:将B状态的循环逻辑放到
b_loop_task函数中,作为独立异步任务启动,不阻塞on_enter_B的执行。 - 任务生命周期管理:在
on_enter_B中启动任务,on_exit_B中取消任务,避免任务泄漏。 - 事件循环调度:循环中的
await asyncio.sleep(1)会主动让出事件循环,让MQTT消息接收任务有机会处理新的状态变更请求。
控制器脚本无需修改
原控制器脚本可以直接使用,发送的to_B、to_A、stop指令现在能正常触发状态机切换。
内容的提问来源于stack exchange,提问作者David Schmider
相关产品推荐
相关产品推荐

