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

Uvicorn运行FastAPI时Paho MQTT publish无消息,如何多线程执行publish?

问题根源

你遇到的MQTT消息发不出去的核心原因是paho-mqtt客户端需要持续运行网络循环处理底层IO、QoS握手逻辑,你只执行了connect和publish操作,没有启动网络循环,消息只是被存入了客户端本地缓冲区,根本没有发送到MQTT broker,所以即使publish返回0也不代表消息成功发出。
另外你代码里存在一个显性错误:初始化mqtt.Client时传入了args['model'],但你定义的args字典中并没有model字段,会导致客户端ID异常,也可能影响连接稳定性。

修复方案

paho-mqtt本身已经提供了loop_start()方法,可以直接在后台启动独立线程处理MQTT网络事件,不需要你手动实现线程逻辑,修改步骤如下:

  1. 修复MQTT客户端初始化参数,删除不存在的args['model']参数,或补充对应配置
  2. 连接MQTT broker后调用loop_start()启动后台网络循环
  3. 可选:添加FastAPI shutdown事件,服务停止时优雅关闭MQTT连接和后台循环
修改后完整代码
import io
import os
import sys
import json
import time
from PIL import Image
import paho.mqtt.client as mqtt
from fastapi import FastAPI, File, HTTPException, UploadFile, Form

# Initialize FastAPI
app = FastAPI()

# Initialize ENV variables
args = {
    'broker'    : os.environ.get('BROKER', '127.0.0.1'),
    'port'      : int(os.environ.get('PORT', '1883')),
    'topic'     : os.environ.get('TOPIC', 'topic'),
    # 可选补充model配置,或者直接删除下面客户端初始化的对应参数
    'model'     : os.environ.get('MODEL', 'fastapi_mqtt_client')
}
    
# Initialize MQTT
print('Connecting to MQTT broker {}:{}.'.format(args['broker'], args['port']), flush=True)
mqtt_client = mqtt.Client(args['model'])
mqtt_client.connect(args['broker'], args['port'])
# 新增:启动MQTT后台网络循环线程
mqtt_client.loop_start()

# 可选:服务停止时优雅关闭MQTT连接
@app.on_event("shutdown")
async def shutdown_event():
    mqtt_client.loop_stop()
    mqtt_client.disconnect()

@app.get("/")
async def info():
    return "Send a POST request to / with an image\nWill publish results to topic {}".format(args['topic'])
    
@app.post("/")
async def run(image: UploadFile = File(...)):
    try:
        start = time.time()
        
        # Read request data
        contents = await image.read()
        image = Image.open(io.BytesIO(contents))

        # Do something with the image
        results = ["todo"]
        print('Process took {} seconds'.format(time.time() - start), flush=True)
        
        # Publish to MQTT topic
        print('Publish to MQTT {}'.format(args['topic']), flush=True)
        (rc, mid) = mqtt_client.publish(args['topic'], json.dumps(results), qos=2)
        print("Code {} while sending message {}: {}".format(rc, mid, mqtt.error_string(rc)))

        # Format response
        data = {}
        data['res'] = results
        data['count'] = len(results)
        data['success'] = True
        return data
    except:
        e = sys.exc_info()[1]
        print('Python error with no Exception handler:')
        print('Traceback error: {}'.format(e))
        raise HTTPException(status_code=500, detail=str(e))
验证说明

修改后重启服务,再发送POST请求,订阅对应MQTT主题即可正常收到消息。如果确实需要自己手动实现线程运行publish逻辑,可以用Python标准库的threading.Thread封装publish操作,启动线程后不阻塞HTTP响应返回,但相比paho自带的loop_start方案需要额外处理线程安全、异常捕获问题,更推荐使用上面的原生方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 23:36:05