You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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())

关键修改点

  1. 补充状态转移:添加to_A和to_B的转移规则,对应控制器发送的指令。
  2. 独立循环任务:将B状态的循环逻辑放到b_loop_task函数中,作为独立异步任务启动,不阻塞on_enter_B的执行。
  3. 任务生命周期管理:在on_enter_B中启动任务,on_exit_B中取消任务,避免任务泄漏。
  4. 事件循环调度:循环中的await asyncio.sleep(1)会主动让出事件循环,让MQTT消息接收任务有机会处理新的状态变更请求。

控制器脚本无需修改

原控制器脚本可以直接使用,发送的to_B、to_A、stop指令现在能正常触发状态机切换。

内容的提问来源于stack exchange,提问作者David Schmider

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 16:39:56