Python MQTT+FastAPI加法服务实现正确性求证
需求理解与代码实现检查
原始需求
- 提供简单REST API接口接收加法输入并返回结果
- 提供MQTT接口能力:
- 订阅
calculator/add主题,接收{"number1":10, "number2":20}格式的加法请求负载 - 计算完成后,向
calculator/add/result主题发布{"final_result":30}格式的结果
- 订阅
你的需求理解偏差
你的理解存在明显偏差,原始需求是两个独立的功能模块:
- REST API模块:直接接收加法输入,计算后返回结果,不需要通过MQTT的发布订阅流程绕圈
- MQTT模块:独立监听
calculator/add主题的加法请求,计算结果后发布到指定结果主题
你错误地把REST API和MQTT的发布订阅绑定成了一个闭环流程,这不符合需求的核心逻辑。
代码实现的核心问题
你的代码多处不符合需求要求,具体问题如下:
MQTT主题配置错误:
- 错误地将发布、订阅主题都设置为
calculator/add/RESULT,但需求要求订阅calculator/add、发布到calculator/add/result - 发布的负载格式错误:代码发布的是
{"a":a,"b":b},但需求要求MQTT模块发布的是{"final_result":计算结果}格式
- 错误地将发布、订阅主题都设置为
MQTT核心逻辑缺失:
- 没有实现订阅
calculator/add主题后处理加法请求、计算结果并发布到结果主题的核心逻辑,当前on_message仅打印收到的内容,完全没有处理计算和结果发布的步骤
- 没有实现订阅
REST接口设计问题:
- 需求未明确要求用GET查询参数,用POST传递请求体(如JSON格式的
number1和number2)是更符合REST规范的设计 - 接口强制限制参数
gt=0,但需求未说明只能计算正数加法,属于多余限制
- 需求未明确要求用GET查询参数,用POST传递请求体(如JSON格式的
其他细节错误:
- FastAPI应用的description字段拼写错误(
resul应为result) - MQTT客户端ID拼写错误(
CALCULATER应为CALCULATOR) - MQTT客户端启动顺序不合理:先启动循环
loop_start()再设置订阅和回调,可能导致连接初期漏收消息
- FastAPI应用的description字段拼写错误(
修正后的代码示例
from fastapi import FastAPI, Query from pydantic import BaseModel import json import paho.mqtt.client as mqtt import logging import os # 日志配置 logging.basicConfig(level=os.environ.get("LOGLEVEL", "INFO"), format='%(levelname)s - %(asctime)s - %(message)s', datefmt='%m-%d %H:%M:%S',) # 初始化FastAPI应用 app = FastAPI(title="My Sweet Calculator", description="A simple service that handles add operations via REST API and MQTT") # MQTT配置常量 BROKER = 'mqtt.eclipseprojects.io' PORT = 1883 MQTT_REQUEST_TOPIC = "calculator/add" MQTT_RESULT_TOPIC = "calculator/add/result" CLIENT_ID = 'CALCULATOR_SERVICE' # MQTT连接成功回调 def on_connect(client, userdata, flags, rc): if rc == 0: logging.info("Connected to MQTT Broker!") # 连接成功后立即订阅加法请求主题 client.subscribe(MQTT_REQUEST_TOPIC) else: logging.error(f"Failed to connect to MQTT Broker, return code {rc}") # MQTT消息处理回调:处理加法请求并发布结果 def on_message(client, userdata, message): try: # 解析收到的MQTT负载 payload = json.loads(message.payload.decode()) number1 = payload.get("number1") number2 = payload.get("number2") # 校验参数合法性 if isinstance(number1, int) and isinstance(number2, int): final_result = number1 + number2 # 发布计算结果到指定主题 client.publish(MQTT_RESULT_TOPIC, json.dumps({"final_result": final_result})) logging.info(f"Processed MQTT add request: {number1} + {number2} = {final_result}, published to {MQTT_RESULT_TOPIC}") else: logging.error("Invalid MQTT payload: must contain integer fields 'number1' and 'number2'") except Exception as e: logging.error(f"Error processing MQTT message: {str(e)}") # 初始化并启动MQTT客户端 mqtt_client = mqtt.Client(CLIENT_ID) mqtt_client.on_connect = on_connect mqtt_client.on_message = on_message mqtt_client.connect(BROKER, PORT) mqtt_client.loop_start() # REST API:POST接口接收加法请求(符合REST规范) class AddRequest(BaseModel): number1: int number2: int @app.post("/add") async def add_post(request: AddRequest): result = request.number1 + request.number2 logging.info(f"Processed REST POST add request: {request.number1} + {request.number2} = {result}") return {"final_result": result} # 可选:支持GET查询参数的加法接口(根据需求选择保留) @app.get("/add") async def add_get(number1: int = Query(...), number2: int = Query(...)): result = number1 + number2 logging.info(f"Processed REST GET add request: {number1} + {number2} = {result}") return {"final_result": result}
内容的提问来源于stack exchange,提问作者Jananath Banuka
相关产品推荐
相关产品推荐

