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

Python MQTT+FastAPI加法服务实现正确性求证

需求理解与代码实现检查

原始需求

  • 提供简单REST API接口接收加法输入并返回结果
  • 提供MQTT接口能力:
    1. 订阅calculator/add主题,接收{"number1":10, "number2":20}格式的加法请求负载
    2. 计算完成后,向calculator/add/result主题发布{"final_result":30}格式的结果

你的需求理解偏差

你的理解存在明显偏差,原始需求是两个独立的功能模块:

  1. REST API模块:直接接收加法输入,计算后返回结果,不需要通过MQTT的发布订阅流程绕圈
  2. MQTT模块:独立监听calculator/add主题的加法请求,计算结果后发布到指定结果主题

你错误地把REST API和MQTT的发布订阅绑定成了一个闭环流程,这不符合需求的核心逻辑。

代码实现的核心问题

你的代码多处不符合需求要求,具体问题如下:

  1. MQTT主题配置错误:

    • 错误地将发布、订阅主题都设置为calculator/add/RESULT,但需求要求订阅calculator/add、发布到calculator/add/result
    • 发布的负载格式错误:代码发布的是{"a":a,"b":b},但需求要求MQTT模块发布的是{"final_result":计算结果}格式
  2. MQTT核心逻辑缺失:

    • 没有实现订阅calculator/add主题后处理加法请求、计算结果并发布到结果主题的核心逻辑,当前on_message仅打印收到的内容,完全没有处理计算和结果发布的步骤
  3. REST接口设计问题:

    • 需求未明确要求用GET查询参数,用POST传递请求体(如JSON格式的number1和number2)是更符合REST规范的设计
    • 接口强制限制参数gt=0,但需求未说明只能计算正数加法,属于多余限制
  4. 其他细节错误:

    • FastAPI应用的description字段拼写错误(resul应为result)
    • MQTT客户端ID拼写错误(CALCULATER应为CALCULATOR)
    • MQTT客户端启动顺序不合理:先启动循环loop_start()再设置订阅和回调,可能导致连接初期漏收消息

修正后的代码示例

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 01:00:56