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

Python实现MQTT订阅消息保存至本地文件的代码故障排查

代码存在的问题
  • on_message是paho-mqtt的事件回调函数,由客户端在收到消息时自动触发执行,不需要手动调用。你代码里注释的手动调用写法既缺少客户端自动传入的三个参数,也无法在消息到达时主动获取数据,完全不符合回调的运行逻辑。
  • 文件写入逻辑完全错误:print('file.txt', file=f)只会向文件写入固定字符串"file.txt",和实际收到的MQTT消息没有任何关系;注释掉的f.write(print(...))写法也不成立,print()函数没有返回值,无法写入有效内容。
  • 数据类型处理错误:msg.payload是字节流(bytes)类型,直接转字符串不做解码会带b''格式标记,不指定编码还容易出现乱码。
  • 全局位置的文件写入逻辑只会在程序启动时执行一次,此时还没有收到任何MQTT消息,根本拿不到有效数据。
  • 追加模式下直接用json.dump写结构化数据,会导致最终文件内容是多个JSON片段拼接的非法格式,后续无法直接用JSON库解析。
正确实现方案

以下是可直接运行的完整实现,基于paho-mqtt库,做了异常兼容和格式规范:

import json
import paho.mqtt.client as mqtt

# 替换为你自己的MQTT服务配置
MQTT_BROKER = "127.0.0.1"
MQTT_PORT = 1883
MQTT_TOPIC = "test/topic"
MQTT_CLIENT_ID = "mqtt_save_file_client"

def on_connect(client, userdata, flags, rc):
    # 连接成功后再订阅,保证断连重连后订阅关系自动恢复
    if rc == 0:
        client.subscribe(MQTT_TOPIC)
    else:
        print(f"MQTT连接失败,错误码:{rc}")

def on_message(client, userdata, msg):
    try:
        # 解码字节流为字符串,默认用utf-8,特殊编码可自行替换
        payload = msg.payload.decode("utf-8")
        # 组装要存储的消息结构
        msg_record = {
            "topic": msg.topic,
            "content": payload,
            "qos": msg.qos,
            "timestamp": msg.timestamp
        }
        # 追加写入文件,指定编码避免乱码
        with open("data.json", "a", encoding="utf-8") as f:
            # 采用JSON Lines格式,每条消息占一行,避免追加破坏JSON结构
            f.write(json.dumps(msg_record, ensure_ascii=False) + "\n")
    except Exception as e:
        print(f"消息写入失败:{str(e)}")

if __name__ == "__main__":
    client = mqtt.Client(client_id=MQTT_CLIENT_ID)
    # 如果你的MQTT服务需要账号密码认证,解开下一行填入对应信息
    # client.username_pw_set("your_username", "your_password")
    client.on_connect = on_connect
    client.on_message = on_message
    client.connect(MQTT_BROKER, MQTT_PORT, keepalive=60)
    # 阻塞运行持续监听消息
    client.loop_forever()
使用说明
  • 把代码里的MQTT连接地址、端口、订阅主题替换成你自己的业务参数即可运行,认证信息按需配置。
  • 存储采用JSON Lines格式,每条消息单独占一行,追加写入不会破坏文件格式,后续处理时逐行加载解析即可。
  • 代码默认用UTF-8编码处理消息和文件,避免Windows系统下默认GBK编码导致的乱码问题。
  • 不要尝试在回调外部定义全局变量暂存消息再统一写入,多线程场景下很容易出现数据竞争丢数据,直接在回调里处理写入是最稳妥的新手友好方案。
  • 如果你的消息吞吐量很高(每秒上百条),可以把文件打开操作移到主程序启动时,把文件对象存在客户端的userdata参数里,回调里直接写入,程序退出前统一关闭文件,减少频繁IO的性能损耗。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 18:39:27