AMQP 协议完整详解
一、AMQP 基础概述
AMQP(Advanced Message Queuing Protocol,高级消息队列协议)
- 定位:面向消息中间件(MQ)的开放标准二进制协议
- 核心目标:实现不同厂商消息中间件互通(互操作性),客户端与Broker无关
- 主流版本:
- AMQP 0-9-1:最广泛使用,RabbitMQ 默认协议(重点讲解)
- AMQP 1.0:全新重构,IBM、ActiveMQ Artemis、Azure Service Bus 支持,与 0-9-1 不兼容
日常开发提到 AMQP,若无特殊说明,一般指 AMQP 0-9-1
核心特性
- 二进制协议:传输效率高于文本协议(HTTP、STOMP)
- 可靠投递模型:支持确认、持久化、事务、死信
- 灵活路由:Exchange + Binding 模型,不止简单队列点对点
- 面向会话,多路复用:一条TCP连接承载多个逻辑通道(Channel)
- 跨平台、跨语言标准
二、整体架构模型(AMQP 0-9-1)
生产者(Publisher) → TCP连接 → Channel → Exchange【交换机】
↓ Binding绑定关系
消费者(Consumer) ← Channel ← Queue【队列】
五大核心组件:Connection、Channel、Exchange、Binding、Queue
1. Connection(TCP连接)
客户端与Broker之间一条TCP物理连接
- 建立连接:握手、协议版本协商、认证(用户名/密码)、虚拟主机(vhost)选定
- 开销较大,应用推荐单进程少量长连接,不要频繁创建销毁
2. Channel(通道)
Connection 内的逻辑轻量级会话
- 一条TCP连接可以创建成千上万Channel,多路复用TCP通道
- AMQP所有命令(发消息、声明队列、消费)都在Channel上执行
- 好处:避免大量TCP连接占用服务器文件句柄;连接复用
重要限制:Channel非线程安全!多线程操作必须各自使用独立Channel**
3. Virtual Host (vhost 虚拟主机)
Broker内部命名隔离空间,类似“命名空间”
- 每个vhost拥有独立Exchange、Queue、权限策略
- RabbitMQ 默认 vhost:/
- 权限控制粒度绑定在 vhost 上
4. Exchange 交换机
消息接收入口,负责路由消息,消息不会直接投递到队列
生产者只发送消息到Exchange,由Exchange根据路由键(routing-key) + 绑定规则分发消息。
AMQP 0-9-1 内置4种标准交换机类型
- Direct(直连交换机)
- 规则:消息
routing-key完全匹配 Binding 的routing-key才投递队列 - 场景:点对点、精准路由、任务分发
- Fanout(扇形交换机)
- 规则:忽略routing-key,消息广播到所有绑定的队列
- 场景:广播通知、日志广播、多消费者同时接收同一条消息
- Topic(主题交换机)
- 规则:支持通配符模糊匹配
*:匹配一个单词#:匹配零个或多个单词
单词之间使用.分割,例:order.create、order.pay.success
- 场景:日志分类、事件订阅(最常用)
- Headers(头交换机)
- 不使用 routing-key,依靠消息
headers属性键值对匹配 - 使用极少,性能较差
tips:RabbitMQ扩展交换机:Dead Letter Exchange(死信交换机DLX)、延迟交换机等,属于厂商扩展,非AMQP标准。
5. Binding 绑定
Exchange 和 Queue 之间的关联关系
- 绑定参数:交换机、队列、binding-key(绑定键)
- 本质:保存路由匹配规则,Broker内部路由查表依据
6. Queue 消息队列
消息最终存储载体,消息存放在队列,等待消费者拉取/推送消费
属性:
- 名称:队列唯一标识
- durable:持久化(Broker重启队列是否保留)
- exclusive:排他队列,仅创建Channel可见,连接断开自动删除
- auto-delete:无消费者连接时,自动删除队列
三、消息结构(AMQP Message)
一条完整消息包含两部分:Header(头部属性) + Payload(二进制消息体)
常用消息属性(Basic.Properties)
content_type:MIME类型(application/json等)
delivery_mode:1=非持久消息,2=持久消息
priority:消息优先级
correlation_id:关联ID(RPC调用常用)
reply_to:RPC回调队列
expiration:消息过期时间(ms)
message_id:消息唯一ID
timestamp:时间戳
type:消息类型标签
user_id / app_id
headers:自定义键值对(headers交换机使用、DLX死信附加信息)
持久化说明
队列durable + delivery_mode=2 同时开启,消息重启不丢失;
仅其中一项生效无法保证持久。
关键字段
- routing-key:消息路由键,发送消息时指定
- exchange:发送目标交换机,空字符串代表默认交换机(Direct类型)
四、AMQP 通信流程(命令模型)
AMQP 采用 RPC命令模型:客户端 ↔ Broker 互相发送方法帧(Method Frame)
所有交互基于帧(Frame),TCP流切分成多种帧类型:
1. Method Frame:核心命令(声明队列、发布消息、消费、ACK)
2. Header Frame:消息属性头
3. Body Frame:消息体(大消息会拆分多个Body帧)
4. Heartbeat Frame:心跳帧,检测连接存活
5. Close Frame:关闭连接/通道
典型完整收发流程
- Client 建立 TCP Connection,握手、认证、选定vhost
- 打开 Channel
- 声明 Exchange(exchange.declare)
- 声明 Queue(queue.declare)
- 绑定 Queue 到 Exchange(queue.bind)
生产者流程:
6. basic.publish 发布消息到Exchange消费者流程二选一:
- 推送模式(主动投递):basic.consume 订阅队列 → Broker主动推送消息 basic.deliver
- 拉取模式:basic.get 单次获取一条消息(不推荐高并发)
- 消费者处理完成,发送 basic.ack 确认;失败可 basic.nack / basic.reject
五、可靠性机制(AMQP标准定义)
1. 生产者确认(Publisher Confirms,RabbitMQ扩展,非原生AMQP 0-9-1标准)
开启 confirm.select
Broker收到消息并持久化完成后,返回 basic.ack;丢失返回 basic.nack
解决问题:生产者不知道Broker是否成功接收消息
2. 事务(AMQP原生 basic.tx)
tx.select → tx.commit / tx.rollback
性能极差,生产环境强烈不推荐,优先使用 Publisher Confirms
3. 消费者消息确认 basic.ack
消费模式分两种:
- 自动确认 auto_ack=true:消息推送后立即标记删除
风险:消费者进程崩溃,消息丢失
- 手动确认 auto_ack=false(推荐)
业务处理成功 → basic.ack(删除消息)
处理失败 → basic.nack
- requeue=true:消息重新放回队列,重新消费
- requeue=false:丢弃消息/进入死信队列
4. 死信(Dead Letter,DLX)
AMQP标准没有DLX,属于RabbitMQ扩展
消息进入死信场景:
1. nack并且requeue=false
2. 消息过期(expiration)
3. 队列达到最大长度被溢出
队列声明时指定参数:
x-dead-letter-exchange、x-dead-letter-routing-key
5. 持久化
- 队列持久化:queue.declare(durable=true)
- 消息持久化:delivery_mode=2
Broker宕机重启后重建队列、加载磁盘消息。
六、流量控制 QoS(basic.qos)
重要:只作用于推送消费模式(basic.consume)
参数 prefetch_count:限制Broker最多推送给同一个消费者未ACK消息数量
作用:防止消费端堆积大量消息引发OOM,实现消费者负载限流
prefetch_count=0:不限制
七、AMQP 0-9-1 VS AMQP 1.0 核心区别
| 维度 | AMQP 0-9-1 | RabbitMQ使用 | AMQP 1.0 |
|---|---|---|---|
| 架构 | Broker内置Exchange/Queue模型,强绑定MQ服务设计 | √ | 通用消息传输协议,不定义队列路由模型,仅定义传输帧 |
| 路由模型 | Exchange+Binding | √ | 协议层无内置路由,交由上层实现 |
| 兼容性 | 和1.0完全不互通 | √ | 全新协议 |
| 代表实现 | RabbitMQ | √ | ActiveMQ Artemis、Azure MQ、IBM MQ |
八、常见开发坑点(基于AMQP协议特性)
- Channel 非线程安全,多线程必须新建独立Channel
- 不要频繁创建Connection,连接很重;Channel可以复用
- auto_ack=true极易丢消息,业务系统禁止使用
- 持久化必须队列+消息delivery_mode同时开启
- basic.get 轮询拉取性能低下,高并发一律用 basic.consume 推送模式
- 心跳(heartbeat)必须合理配置,防止防火墙空闲断开TCP
- 大消息会被切分成多个Body帧,客户端需要重组
- Fanout交换机完全忽略routing-key,代码传值无效
九、AMQP 典型协议交互时序简化
Client ↔ Broker
1. Connection.Start → 协商协议版本、安全机制
2. Connection.Tune → 协商最大帧大小、channel上限
3. Connection.Open(vhost)
4. Channel.Open
====生产者====
Exchange.Declare
Queue.Declare
Queue.Bind
Basic.Publish(message)
====消费者====
Basic.Consume(queue, auto_ack=false)
← Basic.Deliver (Broker推送消息)
业务处理
→ Basic.Ack(delivery_tag)
十、和其他MQ协议简单对比
- AMQP:可靠、标准二进制、适合金融业务,RabbitMQ主力协议
- STOMP:文本协议,简单,功能弱
- MQTT:面向物联网,轻量,不适合业务消息可靠投递
- Kafka协议:Kafka私有协议,高吞吐,没有标准通用规范