如何在Python中无需MQTT客户端解码并区分客户端发送的MQTT消息?
手动解析MQTT协议实现消息解码与类型区分
要解决你的问题,核心是直接解析MQTT协议的报文结构——不需要依赖MQTT客户端库,通过Socket读取字节流后,按照MQTT协议规范手动拆解帧结构即可完成消息解码和类型区分。
核心逻辑:MQTT报文的固定结构
所有MQTT报文都由三部分组成:
- 固定头部:1字节起,包含消息类型、QoS、Dup、Retain等标志,是区分消息类型的关键
- 可变头部:仅部分报文存在,包含Topic、报文ID等元数据
- 负载:实际消息内容(如PUBLISH的消息体)
第一步:通过固定头部区分消息类型
固定头部的第一个字节,高4位是消息类型码,对应MQTT的不同操作:
0x01:CONNECT(客户端连接请求)0x03:PUBLISH(消息发布)0x04:PUBACK(发布确认)0x08:SUBSCRIBE(订阅请求)0x09:SUBACK(订阅确认)0x0B:UNSUBSCRIBE(取消订阅)0x0C:UNSUBACK(取消订阅确认)0x0D:PINGREQ(心跳请求)0x0E:PINGRESP(心跳响应)0x10:DISCONNECT(断开连接)
解析时先读取第一个字节,提取高4位即可判断消息类型:
first_byte = ord(sock.recv(1)) msg_type = (first_byte >> 4) & 0x0F # 提取高4位
第二步:解码客户端的PUBLISH消息
客户端发往Broker的PUBLISH报文是你需要保存的核心,解析流程如下:
- 解析剩余长度:固定头部第二个字节开始是可变长度编码的剩余长度(表示后续报文的总字节数),需要按MQTT规则解析
- 解析Topic:可变头部以2字节的Topic长度开头,后续是UTF-8编码的Topic字符串
- 解析报文ID(可选):如果QoS等级>0,Topic后会跟着2字节的报文ID
- 提取Payload:剩余字节就是实际的消息内容
关键解析函数示例
def parse_mqtt_remaining_length(sock): """解析MQTT的可变长度剩余字段""" remaining_length = 0 multiplier = 1 while True: byte = ord(sock.recv(1)) remaining_length += (byte & 0x7F) * multiplier if not (byte & 0x80): break multiplier *= 128 if multiplier > 128**3: raise ValueError("无效的剩余长度编码") return remaining_length def parse_publish_packet(sock, first_byte, remaining_length): """解析PUBLISH报文,返回包含Topic、Payload的字典""" # 提取QoS等级 qos = (first_byte >> 1) & 0x03 # 解析Topic topic_len_bytes = sock.recv(2) topic_len = (topic_len_bytes[0] << 8) | topic_len_bytes[1] topic = sock.recv(topic_len).decode('utf-8', errors='replace') # 解析报文ID(仅QoS>0时存在) packet_id = None if qos > 0: packet_id_bytes = sock.recv(2) packet_id = (packet_id_bytes[0] << 8) | packet_id_bytes[1] remaining_length -= 2 # 减去报文ID的2字节 # 提取Payload payload_len = remaining_length - 2 - topic_len # 减去Topic长度的2字节和Topic本身 payload = sock.recv(payload_len) return { "topic": topic, "qos": qos, "dup": (first_byte >> 3) & 0x01, "retain": first_byte & 0x01, "packet_id": packet_id, "payload": payload }
第三步:完整的转发与处理流程
在你的Socket转发逻辑中,每次从客户端或Broker读取数据时,先解析报文类型,再针对性处理:
def handle_mqtt_traffic(client_sock, broker_sock): """处理客户端与Broker之间的MQTT流量转发和解析""" while True: # 读取客户端发往Broker的报文 try: first_byte = ord(client_sock.recv(1)) except ConnectionResetError: break msg_type = (first_byte >> 4) & 0x0F remaining_length = parse_mqtt_remaining_length(client_sock) # 如果是PUBLISH消息,解码并保存 if msg_type == 0x03: publish_data = parse_publish_packet(client_sock, first_byte, remaining_length) # 这里添加保存逻辑,比如写入数据库或文件 print(f"已保存PUBLISH消息:Topic={publish_data['topic']},Payload={publish_data['payload']}") # 重新构造完整报文并转发给Broker(需实现剩余长度的反向编码函数) # 示例框架: fixed_header = bytes([first_byte]) # 补充剩余长度的编码字节流 # remaining_len_bytes = encode_remaining_length(remaining_length) # broker_sock.sendall(fixed_header + remaining_len_bytes + ...) else: # 其他类型消息直接转发 remaining_data = client_sock.recv(remaining_length) broker_sock.sendall(bytes([first_byte]) + remaining_data) # 同理处理Broker发往客户端的报文,解析msg_type区分消息类型 # ...
注意事项
- 剩余长度的编码是MQTT解析的核心,必须严格按照协议实现,否则会导致后续字节读取错位
- UTF-8解码时建议添加错误处理(如
errors='replace'),避免因非法字符导致程序崩溃 - 转发报文时必须保证完整转发整个MQTT帧,不能拆分或丢失字节,否则会导致Broker/客户端解析失败
- MQTT 3.1.1和5.0的报文结构略有差异,需根据你使用的协议版本调整解析逻辑
内容的提问来源于stack exchange,提问作者AFORS
相关产品推荐
相关产品推荐

