2.1 消息模型与消息属性


文档摘要

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

2.1 消息模型与消息属性

本节摘要:一条 AMQP 消息由载荷(body)与属性(properties)两部分构成:载荷是业务数据本体,属性是跟着消息走完全程的"随身档案",决定它的持久化、有效期、路由行为与追溯线索。本节把这条消息解剖开,为追踪单后续每一站的判读提供依据。

第 1 章末尾,那条"订单已创建"的消息离开了生产者。在它进入交换机之前,值得先把它拆开看看:你在代码里写下的一个字符串,进入网络后究竟是什么模样。本章六节的追踪都建立在"知道消息身上带什么"之上。

解剖一条消息

AMQP 把消息分成两个部件。载荷(payload)是不透明的字节序列,Broker 从不解读它——你放 JSON、XML、Protobuf 或裸字节都行,消息队列只负责搬运。属性(properties)则是一张结构化的元数据表,Broker 会认真阅读其中若干项,并据此改变消息的命运。

用一张标注图把两者的层次关系钉在墙上:

图 3 消息结构标注:载荷与属性档案

图 2 消息结构标注:载荷与属性档案

属性表里高频使用的字段,按用途归成三类:存活与存储类——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 的报表二进制直接发布,队列瞬间变成存储服务。正确姿势是文件传对象存储,消息里只放下载地址与摘要——消息队列搬的是通知,不是货物本体。

本节要点回顾

  • 两大部件:载荷不透明由业务定义,属性结构化由 Broker 解读;
  • 三类属性:存活存储类决定生死,协作类支撑 RPC 配对,追溯类服务事故排查;
  • 属性全程跟随:从发布到消费逐字保留,是跨服务追踪与协议约定的载体;
  • 两条军规:载荷自描述带版本,追溯字段宁多勿缺。

消息解剖完毕。下一节进入迷宫的第一个分岔口:交换机的四类分拣台,看看同一张消息如何走出四种命运。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U