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
相关产品推荐
相关产品推荐

