Skip to content

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 实例)邮局
ConnectionTCP 连接电话线
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 B

Headers 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: bridge
bash
# 启动集群
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_status

Linux 安装

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 状态
  • 用户权限管理
  • 集群管理

注意事项

  1. 消息可靠性:使用持久化 Exchange、持久化 Queue 和持久化消息确保消息不丢失
  2. 消费者确认:使用 ACK 机制确认消息处理成功,未 ACK 的消息会重新入队
  3. 连接管理:生产环境使用 Connection Pool 复用连接,避免频繁创建销毁
  4. 资源清理:测试环境注意清理未使用的 Exchange 和 Queue
  5. 监控:通过 Web 管理界面或 API 监控队列长度、消息速率等指标

最佳实践

  1. Exchange 类型选择:精确路由用 direct,广播用 fanout,模式匹配用 topic
  2. 队列持久化:生产环境所有队列都设置 durable: true
  3. 消息持久化:关键消息设置 delivery_mode: 2
  4. 消费确认:手动 ACK(auto_ack: false),确保消息处理成功后再确认
  5. 预取计数:设置合理的 prefetch_count,避免消费者负载不均

下一节

继续学习:PHP 客户端操作

参考链接