RabbitMQ 基础
RabbitMQ 是一个开源的消息代理软件(Message Broker),实现了高级消息队列协议(AMQP),为分布式系统提供可靠的消息传递服务。本节将讲解 AMQP 协议核心概念、RabbitMQ 的核心组件(Exchange、Queue、Binding)和消息模型,以及安装部署方法。
前置知识
阅读本节前,建议先了解:网络协议基础 和 PHP 进程与进程间通信
AMQP 协议概述
什么是 AMQP
AMQP(Advanced Message Queuing Protocol,高级消息队列协议)是一个应用层协议规范,定义了消息的格式和消息代理的行为。它是 RabbitMQ 使用的核心协议。
AMQP 协议模型(三层抽象)
+--------------------------------------------------+
| Application Layer(应用层) |
| Producer → Exchange → Binding → Queue → Consumer |
+--------------------------------------------------+
| Session Layer(会话层) |
| Channel(通道复用连接) |
+--------------------------------------------------+
| Transport Layer(传输层) |
| TCP Connection(TCP 连接) |
+--------------------------------------------------+AMQP 核心概念
| 概念 | 说明 | 类比 |
|---|---|---|
| Producer | 消息生产者,发送消息到 Exchange | 寄件人 |
| Consumer | 消息消费者,从 Queue 接收消息 | 收件人 |
| Broker | 消息代理服务器(RabbitMQ 实例) | 邮局 |
| Connection | TCP 连接 | 电话线 |
| Channel | 通道(复用 Connection 的轻量级连接) | 电话分机 |
| Exchange | 交换机(接收消息,按规则路由到 Queue) | 分拣中心 |
| Queue | 队列(存储消息,等待 Consumer 消费) | 信箱 |
| Binding | 绑定(Exchange 到 Queue 的路由规则) | 投递规则 |
| Routing Key | 路由键(消息路由的依据) | 收件地址 |
| vhost | 虚拟主机(逻辑隔离,类似数据库) | 独立邮局 |
Channel 的作用
RabbitMQ 的 TCP 连接创建开销较大,因此使用 Channel 在单个 Connection 上多路复用:
一个 TCP Connection 可以包含多个 Channel
TCP Connection
├── Channel 1: Producer A 发送消息
├── Channel 2: Consumer B 接收消息
├── Channel 3: Consumer C 接收消息
└── Channel 4: 管理 API 调用性能建议
- 一个 Connection 建议不超过 100 个 Channel
- 每个 Consumer 使用独立的 Channel
- 短连接场景(如 HTTP 请求中发消息)可以复用 Connection Pool
核心组件详解
Exchange 类型
Exchange 是消息路由的核心,不同类型决定不同的路由策略。
| Exchange 类型 | 说明 | 路由规则 |
|---|---|---|
direct | 直连交换机 | 精确匹配 Routing Key |
fanout | 扇出交换机 | 忽略 Routing Key,广播到所有绑定的 Queue |
topic | 主题交换机 | 模式匹配 Routing Key(支持 * 和 # 通配符) |
headers | 头部交换机 | 根据消息头属性匹配(性能较差,较少使用) |
Direct Exchange(直连交换机)
Producer → Exchange (direct) → Binding (routing_key="order.created") → Queue A
→ Binding (routing_key="order.cancelled") → Queue B
消息 routing_key="order.created" → 路由到 Queue A
消息 routing_key="order.cancelled" → 路由到 Queue B
消息 routing_key="order.shipped" → 被丢弃(无匹配 Binding)Fanout Exchange(扇出交换机)
Producer → Exchange (fanout) → Queue A → Consumer 1
→ Queue B → Consumer 2
→ Queue C → Consumer 3
所有消息都会广播到所有绑定的 Queue(忽略 Routing Key)
适用于广播通知、日志分发等场景Topic Exchange(主题交换机)
Exchange (topic) → Binding (pattern="order.*.created") → Queue A
→ Binding (pattern="order.#") → Queue B
→ Binding (pattern="*.payment.*") → Queue C
* 匹配一个词
# 匹配零个或多个词
"order.item.created" → Queue A, Queue B
"order.shipped" → Queue B
"user.payment.done" → Queue C
"order.cancelled" → Queue BHeaders Exchange
php
// Headers Exchange 根据消息头匹配(不常用)
$exchange->publish($message, '', [
'headers' => [
'type' => 'notification',
'format' => 'email',
],
]);
$queue->bind($exchange, '', [
'x-match' => 'all', // all=所有 header 都匹配 / any=任一匹配
'type' => 'notification',
]);Queue(队列)
队列是消息的存储容器,Consumer 从队列中获取消息。
php
// 队列属性
$options = [
'passive' => false, // 被动声明(检查队列是否存在,不存在则报错)
'durable' => true, // 持久化(RabbitMQ 重启后队列仍然存在)
'exclusive' => false, // 排他队列(仅当前连接可用,断开后自动删除)
'auto_delete'=> false, // 自动删除(无消费者时自动删除队列)
'arguments' => [
'x-max-priority' => 10, // 优先级队列(0~10)
'x-message-ttl' => 3600000, // 消息过期时间(毫秒)
'x-max-length' => 100000, // 队列最大消息数
'x-overflow' => 'reject-publish', // 超长处理策略
'x-dead-letter-exchange' => 'dlx.exchange', // 死信 Exchange
'x-dead-letter-routing-key' => 'dlx.key', // 死信 Routing Key
],
];Binding(绑定)
Binding 定义了 Exchange 到 Queue 的路由规则:
php
// 基础绑定
$queue->bind($exchangeName, $routingKey);
// 带 header 匹配的绑定
$queue->bind($exchangeName, '', ['headers' => ['type' => 'order']]);消息模型
简单模型(Simple)
一个 Producer,一个 Queue,一个 Consumer:
Producer → Queue → Consumer工作队列模型(Work Queue / Competing Consumer)
一个 Producer,一个 Queue,多个 Consumer 竞争消费:
Producer → Queue → Consumer 1
→ Consumer 2
→ Consumer 3
消息被轮询分发到各 Consumer(Round-robin)
使用 prefetch_count 控制每个 Consumer 未确认的消息数发布/订阅模型(Publish/Subscribe)
一个 Producer,一个 Fanout Exchange,多个 Queue,多个 Consumer:
Producer → Exchange (fanout) → Queue A → Consumer 1
→ Queue B → Consumer 2
所有 Consumer 收到相同的消息副本
适用于广播、通知、日志分发路由模型(Routing)
使用 Direct Exchange,按 Routing Key 精确路由:
Producer → Exchange (direct, routing_key="info") → Queue (info) → Consumer Info
Producer → Exchange (direct, routing_key="error") → Queue (error) → Consumer Error
Producer → Exchange (direct, routing_key="warn") → Queue (warn) → Consumer Warning主题模型(Topics)
使用 Topic Exchange,按模式匹配路由:
Producer → Exchange (topic)
"order.created" → Queue (order.#)
"order.cancelled" → Queue (order.#)
"user.registered" → Queue (user.#)
"user.payment.ok" → Queue (user.*.ok)RPC 模型
使用 Reply Queue 实现 RPC 调用:
Client → Exchange → Queue (rpc_queue) → Server
Server → Exchange → Queue (reply_queue, correlation_id) → Client
Client 发送消息时指定 reply_to(回调队列)和 correlation_id
Server 处理后通过 reply_to 队列发送响应
Client 根据 correlation_id 匹配响应安装与配置
Docker 安装(推荐)
bash
# 单节点 RabbitMQ
docker run -d \
--name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
-e RABBITMQ_DEFAULT_USER=admin \
-e RABBITMQ_DEFAULT_PASS=admin123 \
rabbitmq:3.13-management
# 5672: AMQP 协议端口
# 15672: Web 管理界面端口
# management 标签包含 Web 管理插件Docker Compose 集群
yaml
# docker-compose.yml
version: '3.8'
services:
rabbitmq1:
image: rabbitmq:3.13-management
container_name: rabbitmq1
hostname: rabbitmq1
environment:
- RABBITMQ_ERLANG_COOKIE=rabbitmq_cluster
- RABBITMQ_DEFAULT_USER=admin
- RABBITMQ_DEFAULT_PASS=admin123
ports:
- "5672:5672"
- "15672:15672"
volumes:
- rabbitmq1_data:/var/lib/rabbitmq
networks:
- rabbitmq-net
rabbitmq2:
image: rabbitmq:3.13-management
container_name: rabbitmq2
hostname: rabbitmq2
environment:
- RABBITMQ_ERLANG_COOKIE=rabbitmq_cluster
ports:
- "5673:5672"
- "15673:15672"
volumes:
- rabbitmq2_data:/var/lib/rabbitmq
networks:
- rabbitmq-net
rabbitmq3:
image: rabbitmq:3.13-management
container_name: rabbitmq3
hostname: rabbitmq3
environment:
- RABBITMQ_ERLANG_COOKIE=rabbitmq_cluster
ports:
- "5674:5672"
- "15674:15672"
volumes:
- rabbitmq3_data:/var/lib/rabbitmq
networks:
- rabbitmq-net
volumes:
rabbitmq1_data:
rabbitmq2_data:
rabbitmq3_data:
networks:
rabbitmq-net:
driver: bridgebash
# 启动集群
docker compose up -d
# 集群配置(在 rabbitmq2 上执行)
docker exec rabbitmq2 rabbitmqctl stop_app
docker exec rabbitmq2 rabbitmqctl reset
docker exec rabbitmq2 rabbitmqctl join_cluster rabbitmq@rabbitmq1
docker exec rabbitmq2 rabbitmqctl start_app
# 在 rabbitmq3 上执行
docker exec rabbitmq3 rabbitmqctl stop_app
docker exec rabbitmq3 rabbitmqctl reset
docker exec rabbitmq3 rabbitmqctl join_cluster rabbitmq@rabbitmq1
docker exec rabbitmq3 rabbitmqctl start_app
# 查看集群状态
docker exec rabbitmq1 rabbitmqctl cluster_statusLinux 安装
bash
# 安装 Erlang
curl -s https://packagecloud.io/install/repositories/rabbitmq/erlang/script.rpm.sh | sudo bash
sudo yum install -y erlang
# 安装 RabbitMQ
curl -s https://packagecloud.io/install/repositories/rabbitmq/rabbitmq-server/script.rpm.sh | sudo bash
sudo yum install -y rabbitmq-server
# 启动
sudo systemctl start rabbitmq-server
sudo systemctl enable rabbitmq-server
# 启用管理插件
sudo rabbitmq-plugins enable rabbitmq_management
# 访问 Web 管理界面
# http://server-ip:15672
# 默认用户: guest / guest(仅限 localhost)
# 创建管理员用户
sudo rabbitmqctl add_user admin admin123
sudo rabbitmqctl set_user_tags admin administrator
sudo rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"Web 管理界面
RabbitMQ 自带 Web 管理界面(默认端口 15672),提供以下功能:
- 查看所有 Exchange、Queue、Binding
- 查看消息内容和消息统计
- 手动发送/消费消息(调试)
- 查看 Connection 和 Channel 状态
- 用户权限管理
- 集群管理
注意事项
- 消息可靠性:使用持久化 Exchange、持久化 Queue 和持久化消息确保消息不丢失
- 消费者确认:使用 ACK 机制确认消息处理成功,未 ACK 的消息会重新入队
- 连接管理:生产环境使用 Connection Pool 复用连接,避免频繁创建销毁
- 资源清理:测试环境注意清理未使用的 Exchange 和 Queue
- 监控:通过 Web 管理界面或 API 监控队列长度、消息速率等指标
最佳实践
- Exchange 类型选择:精确路由用 direct,广播用 fanout,模式匹配用 topic
- 队列持久化:生产环境所有队列都设置
durable: true - 消息持久化:关键消息设置
delivery_mode: 2 - 消费确认:手动 ACK(
auto_ack: false),确保消息处理成功后再确认 - 预取计数:设置合理的
prefetch_count,避免消费者负载不均
下一节
继续学习:PHP 客户端操作