【IoT进阶】MQTT + Protobuf 二进制通信实战:从 Python 双向控制到 Wireshark 抓包可视化全解析

摘要:在物联网(IoT)、车联网与工业场景中,相较于冗长的 JSON 文本,MQTT + Google Protobuf 是兼顾高吞吐、极低功耗与强类型扩展的事实工业标准。本文以“智能宠物喝水器”为实战案例,手把手带你完成协议设计、Python 双端通信,并通过 Wireshark + Lua 插件 实现 Protobuf 二进制 Payload 的无缝抓包与字段可视化。


一、 为什么选择 MQTT + Protobuf?

在低功耗传感器(如 NB-IoT 水表、电池供电设备)或大规模车联网中:

  • JSON 文本:键名重复占用大量字节,序列化/反序列化消耗 CPU,流量与电量成本高。
  • Protobuf 二进制:采用 Varint 变长整型压缩与强类型契约,体积比 JSON 缩减 80%~90%,且具备完美的向前/向后兼容性。

二、 协议设计与定义

我们设计一套智能喝水器的物联网通信协议:包含上行遥测状态(Device -> Cloud/App)与下行控制指令(App -> Device)。

1. 协议数据结构设计

① 上行状态报文(Uplink Telemetry)
  • Topic: my_pet_feeder_proto/001/status
 0                   1                   2                   3
 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
|       Device ID (2 Bytes)     | Water Level(1)|  Status (1)   |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
|                      Timestamp (4 Bytes)                      |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
|                 Battery Voltage (Float 4 Bytes)               |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
  • Device ID (uint32): 设备编号
  • Water Level (uint32): 水位 0 ~ 100 (%)
  • Status (Enum): 0=NORMAL, 1=LOW_WATER, 2=EMPTY
  • Timestamp (uint64): Unix 秒级时间戳
  • Battery Voltage (float): 电池实时电压 (V)
② 下行控制报文(Downlink Command)
  • Topic: my_pet_feeder_proto/001/command
 0                   1                   2                   3
 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
|       Device ID (2 Bytes)     |   Action (1)  | Target Val (1)|
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
  • Device ID (uint32): 目标设备编号
  • Action (Enum): 0=REFILL (远程灌装加水), 1=BEEP (蜂鸣寻宠)
  • Target Val (uint32): 加水目标百分比(如 100)

2. Protobuf 契约文件:feeder.proto

新建 feeder.proto:

syntax = "proto3";

package iot.pet;

// 1. 上行状态报文 (设备 -> App)
message FeederStatus {
    uint32 device_id       = 1; // 设备ID
    uint32 water_level     = 2; // 水位百分比 (0-100)
    
    enum StatusEnum {
        NORMAL = 0;             // 正常
        LOW_WATER = 1;          // 低水位预警
        EMPTY = 2;              // 缺水
    }
    StatusEnum status      = 3; // 设备状态枚举
    uint64 timestamp       = 4; // 时间戳
    float battery_voltage  = 5; // 电池电压
}

// 2. 下行控制报文 (App -> 设备)
message FeederCommand {
    uint32 device_id       = 1; // 目标设备ID
    
    enum ActionEnum {
        REFILL = 0;             // 灌装加水
        BEEP = 1;               // 蜂鸣寻宠
    }
    ActionEnum action      = 2; // 指令动作
    uint32 target_val      = 3; // 目标参数值
}
编译生成 Python 类库:
pip install protobuf grpcio-tools paho-mqtt
python -m grpc_tools.protoc -I. --python_out=. feeder.proto

生成 feeder_pb2.py 备用。


三、 Python 双端通信代码实现

1. 设备端(硬件模拟):feeder_device.py

import paho.mqtt.client as mqtt
import time
import threading
import feeder_pb2

TOPIC_STATUS  = "my_pet_feeder_proto/001/status"
TOPIC_COMMAND = "my_pet_feeder_proto/001/command"

DEVICE_ID = 1
water_level = 100
lock = threading.Lock()

def on_message(client, userdata, msg):
    global water_level
    try:
        cmd = feeder_pb2.FeederCommand()
        cmd.ParseFromString(msg.payload) # Protobuf 反序列化
        print(f"\n[收到下行指令] 动作: {feeder_pb2.FeederCommand.ActionEnum.Name(cmd.action)}, 目标值: {cmd.target_val}%")
        
        if cmd.device_id == DEVICE_ID and cmd.action == feeder_pb2.FeederCommand.REFILL:
            with lock:
                water_level = cmd.target_val
            print(f">>> [水泵启动] 灌装加水完成!当前水位: {water_level}% <<<\n")
    except Exception as e:
        print("指令反序列化失败:", e)

def on_connect(client, userdata, flags, rc, properties=None):
    print("[喝水器] 连接 Broker 成功,正在监听控制指令...")
    client.subscribe(TOPIC_COMMAND, qos=1)

client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id="proto_feeder_001")
client.on_connect = on_connect
client.on_message = on_message
client.connect("broker.emqx.io", 1883, keepalive=60)
client.loop_start()

try:
    print("[喝水器] 启动 Protobuf 上报循环...")
    while True:
        with lock:
            if water_level > 0:
                water_level -= 15
                if water_level < 0: water_level = 0
            cur_level = water_level

        # 构造 Protobuf 实体
        status = feeder_pb2.FeederStatus()
        status.device_id = DEVICE_ID
        status.water_level = cur_level
        status.timestamp = int(time.time())
        status.battery_voltage = 3.95
        
        if cur_level > 30:
            status.status = feeder_pb2.FeederStatus.NORMAL
        elif cur_level > 0:
            status.status = feeder_pb2.FeederStatus.LOW_WATER
        else:
            status.status = feeder_pb2.FeederStatus.EMPTY

        # 序列化为二进制流发布
        binary_payload = status.SerializeToString()
        client.publish(TOPIC_STATUS, binary_payload, qos=0)
        print(f"[上报] HEX=[{binary_payload.hex(' ')}] ({len(binary_payload)}B) | 水位={cur_level}% | 电压={status.battery_voltage:.2f}V")
        
        time.sleep(2)
except KeyboardInterrupt:
    client.loop_stop()
    client.disconnect()

2. 控制端(主人 App):owner_app.py

import paho.mqtt.client as mqtt
import time
import feeder_pb2

TOPIC_STATUS  = "my_pet_feeder_proto/001/status"
TOPIC_COMMAND = "my_pet_feeder_proto/001/command"

def on_message(client, userdata, msg):
    try:
        status = feeder_pb2.FeederStatus()
        status.ParseFromString(msg.payload)
        
        level = status.water_level
        status_text = feeder_pb2.FeederStatus.StatusEnum.Name(status.status)
        bar = "█" * (level // 10) + "-" * (10 - level // 10)
        
        print(f"\r[主人 App 看板] 设备#{status.device_id} 水位:[{bar}] {level:3d}% | 状态:{status_text:<9} | 电池:{status.battery_voltage:.2f}V", end="", flush=True)
        if level == 0:
            print("\n⚠️ 小狗没水了!输入 '1' 立即加水: ", end="", flush=True)
    except Exception as e:
        print("\n解析失败:", e)

def on_connect(client, userdata, flags, rc, properties=None):
    print("[主人 App] 连接成功,正在订阅 Protobuf 数据流...")
    client.subscribe(TOPIC_STATUS, qos=0)

client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id="proto_owner_app")
client.on_connect = on_connect
client.on_message = on_message
client.connect("broker.emqx.io", 1883, keepalive=60)
client.loop_start()

print("\n===== 智能喝水器控制台 (输入 '1' 远程加水,'q' 退出) =====")
try:
    while True:
        cmd = input()
        if cmd.strip() in ['1', 'refill']:
            cmd_pb = feeder_pb2.FeederCommand()
            cmd_pb.device_id = 1
            cmd_pb.action = feeder_pb2.FeederCommand.REFILL
            cmd_pb.target_val = 100
            
            client.publish(TOPIC_COMMAND, cmd_pb.SerializeToString(), qos=1)
            print(f"\n---> 已发送 Protobuf 灌装指令!")
        elif cmd.strip().lower() == 'q':
            break
except KeyboardInterrupt:
    pass

client.loop_stop()
client.disconnect()

四、 Wireshark 抓包与 Protobuf 可视化实战

这是整个方案最核心、最惊艳的环节:让 Wireshark 直接将二进制 Payload 还原为结构化字段树。

1. 编写 Wireshark Lua 桥接插件

新建文件 mqtt_feeder_protobuf.lua:

-- MQTT payload decoders for feeder.proto
local protobuf = Dissector.get("protobuf")

-- 注册状态报文解析器
local feeder_status = Proto("feeder_status_pb", "Pet Feeder Status (Protobuf)")
function feeder_status.dissector(tvb, pinfo, tree)
    pinfo.private["pb_msg_type"] = "message,iot.pet.FeederStatus"
    pcall(Dissector.call, protobuf, tvb, pinfo, tree)
end

-- 注册控制指令解析器
local feeder_command = Proto("feeder_command_pb", "Pet Feeder Command (Protobuf)")
function feeder_command.dissector(tvb, pinfo, tree)
    pinfo.private["pb_msg_type"] = "message,iot.pet.FeederCommand"
    pcall(Dissector.call, protobuf, tvb, pinfo, tree)
end

2. Wireshark 详细配置步骤

步骤 ①:配置 .proto 搜索路径
  1. 将 feeder.proto 放置在一个固定目录(例如:C:\Users\xing\PycharmProjects\testMQTT\protobuf-mqtt)。
  2. 打开 Wireshark -> 编辑 (Edit) -> 首选项 (Preferences) -> 展开 Protocols -> Protobuf。
  3. 在 Protobuf search paths 点击 Edit,添加上述目录并勾选启用。
步骤 ②:安装 Lua 插件

将 mqtt_feeder_protobuf.lua 复制到 Wireshark 的插件目录:

  • Windows 路径:C:\Users\<用户名>\AppData\Roaming\Wireshark\plugins\。
步骤 ③:绑定 MQTT Topic 与 Payload 解码器
  1. 打开 Wireshark -> 编辑 (Edit) -> 首选项 (Preferences) -> Protocols -> MQTT。
  2. 找到 Message Decoding 表格,点击 Edit,新增映射:
    • Topic: my_pet_feeder_proto/001/status
    • Payload dissector: feeder_status_pb
    • Decoding: none
  3. 重启 Wireshark 或按下快捷键 Ctrl + Shift + L 重新加载脚本。

五、 抓包验证效果

过滤条件输入:tcp.port == 1883

点击任意一条 Publish Message,可以看到 Wireshark 完美的解析效果:

> MQ Telemetry Transport Protocol, Publish Message
    Topic: my_pet_feeder_proto/001/status
    Message: 0801180220d6b5e6d3062dcdcc7c40
    [Message decoded as: feeder_status_pb]
v Protocol Buffers: iot.pet.FeederStatus
  v Message: iot.pet.FeederStatus
      Field(1): device_id = 1 (uint32)
      Field(3): status = EMPTY(2) (enum)
      Field(4): timestamp = 1786354390 (uint64)
      Field(5): battery_voltage = 3.950000 (float)

原本一长串十六进制字节 0801180220d6...,在 Wireshark 中被清晰地拆解出了 设备 ID、状态枚举、时间戳 以及精准的浮点数 3.95V 电池电压!


六、 总结

通过本方案:

  1. 网络层:利用 MQTT 提供了可靠的 QoS 消息投递与 Pub/Sub 解耦。
  2. 数据层:利用 Protobuf 实现了极度紧凑的二进制序列化,极大降低了 IoT 设备的功耗与流量成本。
  3. 运维/诊断层:利用 Wireshark + Lua 插件打通了从底层抓包到业务字段的可视化,使二进制排错如同看 JSON 一样直观。

该架构正是 Sparkplug B 等现代工业物联网与车联网核心通信方案的标准实践方式。

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐