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

如何在Python中无需MQTT客户端解码并区分客户端发送的MQTT消息?

手动解析MQTT协议实现消息解码与类型区分

要解决你的问题,核心是直接解析MQTT协议的报文结构——不需要依赖MQTT客户端库,通过Socket读取字节流后,按照MQTT协议规范手动拆解帧结构即可完成消息解码和类型区分。

核心逻辑:MQTT报文的固定结构

所有MQTT报文都由三部分组成:

  1. 固定头部:1字节起,包含消息类型、QoS、Dup、Retain等标志,是区分消息类型的关键
  2. 可变头部:仅部分报文存在,包含Topic、报文ID等元数据
  3. 负载:实际消息内容(如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报文是你需要保存的核心,解析流程如下:

  1. 解析剩余长度:固定头部第二个字节开始是可变长度编码的剩余长度(表示后续报文的总字节数),需要按MQTT规则解析
  2. 解析Topic:可变头部以2字节的Topic长度开头,后续是UTF-8编码的Topic字符串
  3. 解析报文ID(可选):如果QoS等级>0,Topic后会跟着2字节的报文ID
  4. 提取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 02:51:06