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

MQTT接收JSON消息触发JsondecodeError line1 column1错误的问题咨询

错误触发原因

该报错的核心原因是json.loads()方法接收到了空字符串输入,具体触发场景对应如下:

  • MQTT协议为了维持长连接,broker会在长时间无业务消息往来时,按保活周期发送心跳报文,这类报文的payload为空,会被投递到on_message回调函数中直接参与解析,触发JSON解析错误。
  • 网络波动、客户端重连时也可能收到空payload的异常消息,现有代码没有做任何前置校验,直接解析就会崩溃。
  • 此外代码中未做异常捕获,只要收到的消息不是UTF-8编码、不是合法JSON格式、或者缺失employees字段,都会直接抛出异常终止运行。
修复方案

核心修改点

  • 新增payload非空校验,过滤空消息
  • 捕获解码、JSON解析阶段的异常,避免单次异常导致整个程序崩溃
  • 新增字段合法性校验,确保后续DataFrame转换逻辑可正常执行
  • 统一文件路径的写法,避免相对路径和绝对路径混用导致的header重复写入问题
  • 显式设置MQTT保活周期,减少心跳报文的触发频率
  • 新增自动重连配置,避免断网后程序无法自动恢复

修改后的完整可运行代码

import paho.mqtt.client as mqtt
import pandas as pd
import json
import datetime
import os

def on_connect(client,userdata,flag,rc):
    print("Connect",str(rc))
    print("client",client)
    print("flag= ",flag)

def on_message(client,userdata,message):
    print(datetime.datetime.now())
    # 过滤空payload消息
    if not message.payload:
        print("收到空payload消息,跳过处理")
        return
    try:
        # 解码、解析JSON统一加异常捕获
        data = message.payload.decode("utf-8")
        json_data = json.loads(data)
        # 校验必要字段存在且格式符合要求
        if "employees" not in json_data or not isinstance(json_data["employees"], list) or len(json_data["employees"]) == 0:
            print("JSON数据不符合预期格式,跳过处理")
            return
        info = json_data["employees"]
    except (UnicodeDecodeError, json.JSONDecodeError) as e:
        print(f"消息处理失败,错误信息:{str(e)},跳过处理")
        return

    keylist = list(info[0].keys())
    df = pd.DataFrame(info, columns = keylist)
    df.insert(0,"received time",datetime.datetime.now())
  
    # 统一使用绝对路径
    file_path = '/home/code/jsontocsvtest.csv'
    if os.path.isfile(file_path):
        df.to_csv(file_path, index = False, mode='a', header = False)
    else:
        df.to_csv(file_path, index = False, mode='a')

broker_address="0.0.0.0"
client1 = mqtt.Client("client1")
# 设置自动重连,避免断网后程序无法恢复
client1.reconnect_delay_set(min_delay=1, max_delay=120)
client1.on_connect = on_connect
# 显式设置保活周期为300秒(5分钟)
client1.connect(broker_address, keepalive=300)
client1.subscribe("outTopic") 
client1.on_message = on_message
client1.loop_forever()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 04:24:00