【IoT进阶】MQTT + Protobuf 二进制通信实战:从 Python 双向控制到 Wireshark 抓包可视化全解析
·
【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 搜索路径
- 将
feeder.proto放置在一个固定目录(例如:C:\Users\xing\PycharmProjects\testMQTT\protobuf-mqtt)。 - 打开 Wireshark -> 编辑 (Edit) -> 首选项 (Preferences) -> 展开 Protocols -> Protobuf。
- 在 Protobuf search paths 点击 Edit,添加上述目录并勾选启用。
步骤 ②:安装 Lua 插件
将 mqtt_feeder_protobuf.lua 复制到 Wireshark 的插件目录:
- Windows 路径:
C:\Users\<用户名>\AppData\Roaming\Wireshark\plugins\。
步骤 ③:绑定 MQTT Topic 与 Payload 解码器
- 打开 Wireshark -> 编辑 (Edit) -> 首选项 (Preferences) -> Protocols -> MQTT。
- 找到 Message Decoding 表格,点击 Edit,新增映射:
- Topic:
my_pet_feeder_proto/001/status - Payload dissector:
feeder_status_pb - Decoding:
none
- Topic:
- 重启 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 电池电压!
六、 总结
通过本方案:
- 网络层:利用 MQTT 提供了可靠的 QoS 消息投递与 Pub/Sub 解耦。
- 数据层:利用 Protobuf 实现了极度紧凑的二进制序列化,极大降低了 IoT 设备的功耗与流量成本。
- 运维/诊断层:利用 Wireshark + Lua 插件打通了从底层抓包到业务字段的可视化,使二进制排错如同看 JSON 一样直观。
该架构正是 Sparkplug B 等现代工业物联网与车联网核心通信方案的标准实践方式。
更多推荐




所有评论(0)