如何用FastAPI+Paho-MQTT实现后台MQTT消息接收且不阻塞API启动?
解决FastAPI与Paho-MQTT后台订阅阻塞问题
问题根源
你的核心问题在于Paho-MQTT的loop_forever()是同步阻塞方法,无论放在async任务还是主线程中,都会占用当前线程执行权,导致FastAPI无法正常启动:
- 用
asyncio.create_task启动时,同步阻塞的loop_forever()会卡住asyncio事件循环,使FastAPI的startup流程无法完成。 - 在
uvicorn.run前调用时,loop_forever()直接阻塞主线程,根本不会执行到API启动代码。
另外,on_message回调是同步执行的,直接调用Motor异步客户端的insert_one会引发线程安全问题,需要特殊处理。
解决方案:将MQTT订阅放在独立线程运行
把MQTT的循环逻辑放到单独的守护线程中,既不阻塞FastAPI的事件循环,也能持续接收MQTT消息。同时通过asyncio.run_coroutine_threadsafe将MongoDB的异步插入任务提交到FastAPI的事件循环中,保证异步操作的线程安全。
1. 修改mqtt.py代码
import json import asyncio import paho.mqtt.client as mqtt_client import motor.motor_asyncio # 你的配置常量 URL = "mongodb://your-mongo-url" CLIENT_ID = "your-client-id" BROKER = "your-broker-url" PORT = 8883 ROOT_CA_PATH = "path/to/rootCA.pem" USERNAME = "your-mqtt-username" PASSWORD = "your-mqtt-password" TOPIC = "your/topic" # 初始化MongoDB异步客户端 client = motor.motor_asyncio.AsyncIOMotorClient(URL) db = client.collection def connect_mqtt() -> mqtt_client: def on_connect(client, userdata, flags, rc): if rc == 0: print("Connected to MQTT Broker!") else: print(f"Failed to connect, return code {rc}") client = mqtt_client.Client(CLIENT_ID) client.tls_set(ca_certs=ROOT_CA_PATH) client.username_pw_set(USERNAME, PASSWORD) client.on_connect = on_connect client.connect(BROKER, PORT) return client def subscribe(client: mqtt_client): # 获取FastAPI的asyncio事件循环 loop = asyncio.get_event_loop() def on_message(client, userdata, msg): try: payload = msg.payload.decode('utf-8') data = json.loads(payload) # 将异步插入任务提交到事件循环 future = asyncio.run_coroutine_threadsafe(db["devices"].insert_one(data), loop) # 可选:获取插入结果(如果需要处理成功/失败) result = future.result() print(f"Document inserted! ID: {result.inserted_id}") except Exception as e: print(f"Error processing message: {str(e)}") client.subscribe(TOPIC, qos=0) client.on_message = on_message def mqtt_subscribe(): print("Starting MQTT subscriber...") client = connect_mqtt() subscribe(client) client.loop_forever()
2. 修改app.py的启动逻辑
在startup事件中启动MQTT订阅线程,设置为守护线程(随主线程退出而终止):
from fastapi import FastAPI import threading from mqtt import mqtt_subscribe app = FastAPI() @app.on_event("startup") def startup_event(): # 启动MQTT订阅守护线程 mqtt_thread = threading.Thread(target=mqtt_subscribe, daemon=True) mqtt_thread.start() @app.get("/") async def root(): return {"hello": "world"} if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)
验证效果
启动API后,你会看到类似以下输出:
INFO: Started server process [xxxx] INFO: Waiting for application startup. Starting MQTT subscriber... Connected to MQTT Broker! INFO: Application startup complete. INFO: Uvicorn running on http://0.0.0.0:8000 (Press CTRL+C to quit)
此时API路由可以正常访问,后台也能持续接收MQTT消息并存入MongoDB。
内容的提问来源于stack exchange,提问作者Isaquehg
相关产品推荐
相关产品推荐

