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网络事件,不需要你手动实现线程逻辑,修改步骤如下:
- 修复MQTT客户端初始化参数,删除不存在的
args['model']参数,或补充对应配置 - 连接MQTT broker后调用
loop_start()启动后台网络循环 - 可选:添加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
相关产品推荐
相关产品推荐

