2.1 消息模型与消息属性 本节摘要:一条 AMQP 消息由载荷(body)与属性(properties)两部分构成:载荷是业务数据本体,属性是跟着消息走完全程的"随身档案",决定它的持久化、有效期、路由行为与追溯线索。本节把这条消息解剖开,为追踪单后续每一站的判读提供依据。 第 1 章末尾,那条"订单已创建"的消息离开了生产者。在它进入交换机之前,值得先把它拆开看看:你在代码里写下的一个字符串,进入网络后究竟是什么模样。本章六节的追踪都建立在"知道消息身上带什么"之上。 解剖一条消息 AMQP 把消息分成两个部件。载荷(payload)是不透明的字节序列,Broker 从不解读它——你放 JSON、XML、Protobuf 或裸字节都行,消息队列只负责搬运。
本节摘要:一条 AMQP 消息由载荷(body)与属性(properties)两部分构成:载荷是业务数据本体,属性是跟着消息走完全程的"随身档案",决定它的持久化、有效期、路由行为与追溯线索。本节把这条消息解剖开,为追踪单后续每一站的判读提供依据。
第 1 章末尾,那条"订单已创建"的消息离开了生产者。在它进入交换机之前,值得先把它拆开看看:你在代码里写下的一个字符串,进入网络后究竟是什么模样。本章六节的追踪都建立在"知道消息身上带什么"之上。
AMQP 把消息分成两个部件。载荷(payload)是不透明的字节序列,Broker 从不解读它——你放 JSON、XML、Protobuf 或裸字节都行,消息队列只负责搬运。属性(properties)则是一张结构化的元数据表,Broker 会认真阅读其中若干项,并据此改变消息的命运。
用一张标注图把两者的层次关系钉在墙上:

属性表里高频使用的字段,按用途归成三类:存活与存储类——delivery_mode(2 表示持久化)、expiration(单条消息的过期毫秒数)、priority(0 到 9 的优先级);协作类——correlation_id 与 reply_to 是 RPC 模式的搭档,请求消息带上编号,响应消息原样带回,消费方靠它对上号;追溯类——message_id、timestamp、app_id、headers 自定义键值对,排查事故时这些是唯一的案发现场证物。
背景:延续 1.3 节的环境,现在发两条内容相同、属性不同的消息,观察它们在队列中的不同待遇。
操作:
import pika, json connection = pika.BlockingConnection(pika.ConnectionParameters(host="localhost")) channel = connection.channel() event = json.dumps({"order_id": "A1025", "event": "created"}).encode() # 消息一:完整属性档案,持久化、带追溯信息 channel.basic_publish( exchange="trace.direct", routing_key="order.created", body=event, properties=pika.BasicProperties( content_type="application/json", delivery_mode=2, message_id="evt-a1025-created-001", app_id="order-service", timestamp=1796026800, headers={"trace_id": "t-8848"})) # 消息二:裸奔,无任何属性 channel.basic_publish( exchange="trace.direct", routing_key="order.created", body=event) print("两条消息均已投递")
再用命令行验证持久化标记是否生效(persistent 列只在队列与消息都持久化时计数):
rabbitmqctl list_queues name messages persistent # 预期输出: # Listing queues ... # trace.orders 2 1
结果:两条消息都在队列里,但 persistent 列只有 1——带 delivery_mode=2 的那条落了盘,裸奔那条只存在内存。解读:此刻重启 Broker 再查队列,会发现只剩 1 条消息。属性不是可有可无的注解,而是直接写进消息生死簿的判词。追踪单风险点三号(属性缺失导致重启丢失)在此确认。
变式一:验证 expiration。给消息加 expiration="3000"(三秒过期),发布后等待四秒再查队列,消息数归零——过期消息被静默移除,如果队列配了死信交换机则转入死信,这是第 3 章延迟方案的原理地基。变式二:验证属性在消费端的完整还原:
def on_message(ch, method, properties, body): # properties 对象携带发布时设置的全部属性,逐字不动 print("message_id:", properties.message_id) # evt-a1025-created-001 print("trace_id:", properties.headers.get("trace_id")) # t-8848 channel.basic_consume(queue="trace.orders", on_message_callback=on_message) channel.start_consuming() # 运行输出: # message_id: evt-a1025-created-001 # trace_id: t-8848
属性一路无损跟随消息走完发布、路由、存储、投递全程——这就是"随身档案"的含义,也是分布式追踪能跨服务串联的前提。
第一条:载荷只放事实,格式自描述。content_type 必须设置;消息体用团队约定的统一格式(JSON 是常见起点),并带 schema 版本号字段,消费方遇到不认识的版本能显式拒绝而不是解析出垃圾数据。第二条:追溯字段宁多勿缺。trace_id、message_id、时间戳在平日是冗余,在事故夜是救命稻草。多传几个字段的开销是几十字节,少传它们的代价可能是彻夜排查。
还有一个升级期才显出价值的习惯:载荷里带上"产生时间与产生实例"。发布时刻的时间戳与主机标识,是判断"这条消息在路上走了多久"的唯一依据——消费延迟的定位、堆积的时间还原、重复投递的判别,全都依赖它。字段只要两个字,价值贯穿整条链路。
⚠️ 常见坑:把大文件塞进载荷。有人把十几 MB 的报表二进制直接发布,队列瞬间变成存储服务。正确姿势是文件传对象存储,消息里只放下载地址与摘要——消息队列搬的是通知,不是货物本体。
消息解剖完毕。下一节进入迷宫的第一个分岔口:交换机的四类分拣台,看看同一张消息如何走出四种命运。