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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 05:07:44