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

核心特性

  1. 二进制协议:传输效率高于文本协议(HTTP、STOMP)
  2. 可靠投递模型:支持确认、持久化、事务、死信
  3. 灵活路由:Exchange + Binding 模型,不止简单队列点对点
  4. 面向会话,多路复用:一条TCP连接承载多个逻辑通道(Channel)
  5. 跨平台、跨语言标准

二、整体架构模型(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种标准交换机类型

  1. Direct(直连交换机)
  • 规则:消息 routing-key 完全匹配 Binding 的 routing-key 才投递队列
  • 场景:点对点、精准路由、任务分发
  1. Fanout(扇形交换机)
  • 规则:忽略routing-key,消息广播到所有绑定的队列
  • 场景:广播通知、日志广播、多消费者同时接收同一条消息
  1. Topic(主题交换机)
  • 规则:支持通配符模糊匹配
    • *:匹配一个单词
    • #:匹配零个或多个单词
      单词之间使用 . 分割,例:order.createorder.pay.success
  • 场景:日志分类、事件订阅(最常用)
  1. 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:关闭连接/通道

典型完整收发流程

  1. Client 建立 TCP Connection,握手、认证、选定vhost
  2. 打开 Channel
  3. 声明 Exchange(exchange.declare)
  4. 声明 Queue(queue.declare)
  5. 绑定 Queue 到 Exchange(queue.bind)

生产者流程:
6. basic.publish 发布消息到Exchange

消费者流程二选一:
- 推送模式(主动投递):basic.consume 订阅队列 → Broker主动推送消息 basic.deliver
- 拉取模式:basic.get 单次获取一条消息(不推荐高并发)

  1. 消费者处理完成,发送 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-exchangex-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-1RabbitMQ使用AMQP 1.0
架构Broker内置Exchange/Queue模型,强绑定MQ服务设计通用消息传输协议,不定义队列路由模型,仅定义传输帧
路由模型Exchange+Binding协议层无内置路由,交由上层实现
兼容性和1.0完全不互通全新协议
代表实现RabbitMQActiveMQ Artemis、Azure MQ、IBM MQ

八、常见开发坑点(基于AMQP协议特性)

  1. Channel 非线程安全,多线程必须新建独立Channel
  2. 不要频繁创建Connection,连接很重;Channel可以复用
  3. auto_ack=true极易丢消息,业务系统禁止使用
  4. 持久化必须队列+消息delivery_mode同时开启
  5. basic.get 轮询拉取性能低下,高并发一律用 basic.consume 推送模式
  6. 心跳(heartbeat)必须合理配置,防止防火墙空闲断开TCP
  7. 大消息会被切分成多个Body帧,客户端需要重组
  8. 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私有协议,高吞吐,没有标准通用规范