> ## Content Index
> Fetch the complete content index at: https://blog.vercanti.com/llms.txt
> Use this file to discover other available public pages before exploring further.

# RabbitMQ 完全指南
- URL: https://blog.vercanti.com/rabbitmq-wan-quan-zhi-nan/
- Published: 2026-08-28T14:35:42.000Z
- Updated: 2026-08-28T14:59:23.000Z
- Description: RabbitMQ 是基于 AMQP 协议的消息代理，以灵活的路由能力和低延迟著称。 pika.ConnectionParameters 参数： channel.exchange_declare 参数： channel.queue_declare 参数： channel.basic_publish 参数： pika.BasicProperties 常用字段： channel.basic_consume 参数： 防止某个消费者堆积过多未处理消息： aio_pika.connect_robust 参数： connect_robust 与 connect 的区别
- Author: yellowdog
- Tags: 消息队列

> 官方文档：<https://www.rabbitmq.com/documentation.html>  
> 适用版本：RabbitMQ 3.12+（2026-05-07 整理）

RabbitMQ 是基于 AMQP 协议的消息代理，以灵活的路由能力和低延迟著称。

---

## 核心概念

### 组件关系

```
Producer
    |
    | 发布消息（指定 exchange + routing_key）
    v
Exchange（交换机）
    |
    | 按 Binding 规则将消息路由到一个或多个队列
    v
Queue（队列）
    |
    | 按顺序投递消息
    v
Consumer（消费者）

```

- **Producer**：消息生产者，将消息发送到 Exchange，不直接操作 Queue。
- **Exchange**：接收 Producer 的消息，根据路由规则分发到绑定的 Queue。
- **Queue**：存储消息，消费者从此处取消息。消费后消息默认删除。
- **Binding**：Exchange 和 Queue 之间的绑定关系，附带 binding\_key 用于路由匹配。
- **Consumer**：订阅 Queue，接收并处理消息。
- **Message**：由 payload（消息体）和 properties（元数据）组成。

### Exchange 类型对比

| Exchange 类型 | 路由规则                                                   | 典型场景                |
| ----------- | ------------------------------------------------------ | ------------------- |
| direct      | routing\_key 与 binding\_key 完全匹配                       | 任务分发，将不同类型任务路由到对应队列 |
| topic       | routing\_key 与 binding\_key 模式匹配，\* 匹配一个单词，# 匹配零个或多个单词 | 日志系统，按模块和级别组合路由     |
| fanout      | 忽略 routing\_key，广播到所有绑定的 Queue                         | 广播通知，发布/订阅模式        |
| headers     | 根据消息 headers 属性匹配，忽略 routing\_key                      | 按消息属性路由，灵活但性能较差     |

---

## Python 客户端：pika（同步）

### 安装

```bash
pip install pika

```

### 建立连接

```python
import pika

# 基础连接
connection = pika.BlockingConnection(
    pika.ConnectionParameters(host="localhost")
)
channel = connection.channel()

```

`pika.ConnectionParameters` 参数：

| 参数                           | 类型               | 默认值         | 说明                                         |
| ---------------------------- | ---------------- | ----------- | ------------------------------------------ |
| host                         | str              | "localhost" | RabbitMQ 服务器地址                             |
| port                         | int              | 5672        | AMQP 端口                                    |
| virtual\_host                | str              | "/"         | 虚拟主机，用于逻辑隔离                                |
| credentials                  | PlainCredentials | guest/guest | 认证凭据，用 pika.PlainCredentials(user, pwd) 构造 |
| heartbeat                    | int              | 60          | 心跳间隔（秒），0 表示禁用                             |
| connection\_attempts         | int              | 1           | 连接重试次数                                     |
| retry\_delay                 | float            | 2.0         | 重试间隔（秒）                                    |
| socket\_timeout              | float            | 10.0        | Socket 连接超时（秒）                             |
| blocked\_connection\_timeout | float            | None        | 连接被 Broker 阻塞的超时时间                         |

### 声明 Exchange

```python
channel.exchange_declare(
    exchange="logs",
    exchange_type="topic",
    durable=True,
    auto_delete=False,
)

```

`channel.exchange_declare` 参数：

| 参数             | 类型   | 默认值      | 说明                                   |
| -------------- | ---- | -------- | ------------------------------------ |
| exchange       | str  | 必填       | Exchange 名称                          |
| exchange\_type | str  | "direct" | 类型：direct / topic / fanout / headers |
| durable        | bool | False    | 是否持久化，Broker 重启后 Exchange 仍然存在       |
| auto\_delete   | bool | False    | 所有绑定解除后自动删除                          |
| passive        | bool | False    | 仅检查 Exchange 是否存在，不创建                |
| arguments      | dict | None     | 额外参数                                 |

### 声明 Queue

```python
result = channel.queue_declare(
    queue="task_queue",
    durable=True,
    exclusive=False,
    auto_delete=False,
    arguments={"x-max-length": 10000},
)
print(result.method.queue)  # 获取队列名（匿名队列时有用）

```

`channel.queue_declare` 参数：

| 参数           | 类型   | 默认值   | 说明                                          |
| ------------ | ---- | ----- | ------------------------------------------- |
| queue        | str  | ""    | 队列名，空字符串由 Broker 生成随机名                      |
| durable      | bool | False | 持久化队列，Broker 重启后队列不丢失                       |
| exclusive    | bool | False | 独占队列，仅当前连接可用，连接断开后自动删除                      |
| auto\_delete | bool | False | 最后一个消费者断开后自动删除                              |
| passive      | bool | False | 仅检查队列是否存在                                   |
| arguments    | dict | None  | 额外参数，如 x-message-ttl、x-dead-letter-exchange |

### 发布消息

```python
import pika

channel.basic_publish(
    exchange="logs",
    routing_key="error.critical",
    body=b"something went wrong",
    properties=pika.BasicProperties(
        delivery_mode=2,        # 消息持久化
        content_type="text/plain",
        expiration="60000",     # 消息 TTL，单位毫秒
    ),
)

```

`channel.basic_publish` 参数：

| 参数           | 类型              | 默认值   | 说明                                                      |
| ------------ | --------------- | ----- | ------------------------------------------------------- |
| exchange     | str             | 必填    | 目标 Exchange，空字符串表示默认 Exchange（direct，routing\_key 即队列名） |
| routing\_key | str             | 必填    | 路由键                                                     |
| body         | bytes           | 必填    | 消息体，必须为 bytes                                           |
| properties   | BasicProperties | None  | 消息属性，见下方说明                                              |
| mandatory    | bool            | False | 若消息无法路由到任何队列则返回给 Producer                               |

`pika.BasicProperties` 常用字段：

| 字段                | 说明                              |
| ----------------- | ------------------------------- |
| delivery\_mode    | 1 非持久化，2 持久化（需配合 durable queue） |
| content\_type     | 消息体的 MIME 类型                    |
| content\_encoding | 编码方式，如 "utf-8"                  |
| headers           | 自定义消息头（dict）                    |
| expiration        | 消息过期时间（毫秒字符串），超时进入死信队列          |
| message\_id       | 消息唯一标识                          |
| correlation\_id   | 用于 RPC 场景，关联请求与响应               |
| reply\_to         | RPC 回调队列名                       |
| priority          | 消息优先级（0-9），需队列支持 x-max-priority |

### 消费消息

```python
def on_message(channel, method, properties, body):
    print(f"received: {body.decode()}")
    channel.basic_ack(delivery_tag=method.delivery_tag)  # 手动确认

channel.basic_consume(
    queue="task_queue",
    on_message_callback=on_message,
    auto_ack=False,
)
channel.start_consuming()

```

`channel.basic_consume` 参数：

| 参数                    | 类型       | 默认值   | 说明                                               |
| --------------------- | -------- | ----- | ------------------------------------------------ |
| queue                 | str      | 必填    | 消费的队列名                                           |
| on\_message\_callback | callable | 必填    | 收到消息时的回调，签名为 (channel, method, properties, body) |
| auto\_ack             | bool     | False | 自动确认，消息投递后立即 ack，不等待处理完成                         |
| exclusive             | bool     | False | 独占消费，不允许其他消费者消费该队列                               |
| consumer\_tag         | str      | None  | 消费者标识，不填由 Broker 生成                              |

### 手动 ACK

```python
# 确认单条消息
channel.basic_ack(delivery_tag=method.delivery_tag)

# 拒绝并重新入队
channel.basic_nack(delivery_tag=method.delivery_tag, requeue=True)

# 拒绝并丢弃（或进入死信队列）
channel.basic_reject(delivery_tag=method.delivery_tag, requeue=False)

```

### 公平分发

防止某个消费者堆积过多未处理消息：

```python
# 每次只给消费者推送 1 条未 ack 的消息
channel.basic_qos(prefetch_count=1)

```

---

## Python 客户端：aio-pika（异步）

### 安装

```bash
pip install aio-pika

```

### 异步连接

```python
import aio_pika

connection = await aio_pika.connect_robust(
    "amqp://guest:guest@localhost/",
    heartbeat=60,
)

```

`aio_pika.connect_robust` 参数：

| 参数                  | 类型    | 默认值         | 说明                                           |
| ------------------- | ----- | ----------- | -------------------------------------------- |
| url                 | str   | 必填          | AMQP URL，格式 amqp://user:pass@host:port/vhost |
| host                | str   | "localhost" | 与 url 二选一                                    |
| port                | int   | 5672        | AMQP 端口                                      |
| login               | str   | "guest"     | 用户名                                          |
| password            | str   | "guest"     | 密码                                           |
| virtualhost         | str   | "/"         | 虚拟主机                                         |
| heartbeat           | int   | 60          | 心跳间隔（秒）                                      |
| reconnect\_interval | float | 5.0         | 断线重连间隔（秒），connect\_robust 专有                 |

`connect_robust` 与 `connect` 的区别：`connect_robust` 在连接断开后会自动重连，生产环境应始终使用 `connect_robust`。

### 异步发布消息

```python
import aio_pika

async def publish():
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    async with connection:
        channel = await connection.channel()
        exchange = await channel.declare_exchange(
            "direct_logs",
            aio_pika.ExchangeType.DIRECT,
            durable=True,
        )
        await exchange.publish(
            aio_pika.Message(
                body=b"hello world",
                delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
            ),
            routing_key="info",
        )

```

### 异步消费消息

```python
import aio_pika

async def consume():
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    async with connection:
        channel = await connection.channel()
        await channel.set_qos(prefetch_count=10)

        queue = await channel.declare_queue("task_queue", durable=True)

        async with queue.iterator() as queue_iter:
            async for message in queue_iter:
                async with message.process():
                    # process() 上下文管理器自动 ack，异常时自动 nack
                    print(message.body.decode())

```

### 与 FastAPI 集成

```python
from contextlib import asynccontextmanager
import aio_pika
from fastapi import FastAPI

rabbitmq_connection = None

@asynccontextmanager
async def lifespan(app: FastAPI):
    global rabbitmq_connection
    rabbitmq_connection = await aio_pika.connect_robust(
        "amqp://guest:guest@localhost/"
    )
    yield
    await rabbitmq_connection.close()

app = FastAPI(lifespan=lifespan)

@app.post("/publish")
async def publish_message(body: str):
    channel = await rabbitmq_connection.channel()
    exchange = await channel.get_exchange("my_exchange")
    await exchange.publish(
        aio_pika.Message(body=body.encode()),
        routing_key="tasks",
    )
    return {"status": "published"}

```

---

## 高级特性

### 消息持久化

要确保消息在 Broker 重启后不丢失，需同时满足三个条件：

```python
# 1. 声明持久化 Exchange
channel.exchange_declare(exchange="my_ex", exchange_type="direct", durable=True)

# 2. 声明持久化 Queue
channel.queue_declare(queue="my_queue", durable=True)

# 3. 发布时设置 delivery_mode=2
channel.basic_publish(
    exchange="my_ex",
    routing_key="key",
    body=b"data",
    properties=pika.BasicProperties(delivery_mode=2),
)

```

三个条件缺一不可，任意一处未配置持久化均可能丢失消息。

### 死信队列（DLX）

死信队列（Dead Letter Exchange）接收无法被正常消费的消息。消息变为死信的条件：

- 消费者调用 `basic_reject` 或 `basic_nack` 且 `requeue=False`
- 消息 TTL 过期
- 队列达到最大长度

```python
# 1. 声明死信 Exchange 和队列
channel.exchange_declare(exchange="dlx", exchange_type="direct", durable=True)
channel.queue_declare(queue="dead_letter_queue", durable=True)
channel.queue_bind(queue="dead_letter_queue", exchange="dlx", routing_key="dead")

# 2. 声明业务队列时绑定死信配置
channel.queue_declare(
    queue="business_queue",
    durable=True,
    arguments={
        "x-dead-letter-exchange": "dlx",       # 死信发往的 Exchange
        "x-dead-letter-routing-key": "dead",   # 死信的 routing_key
        "x-message-ttl": 30000,                # 队列内消息的过期时间（毫秒）
        "x-max-length": 10000,                 # 队列最大消息数
    },
)

```

### 延迟队列（TTL + DLX 实现）

RabbitMQ 原生不支持延迟投递，通过"TTL 队列 + 死信队列"模拟：

```python
# 延迟队列：消息在此队列中等待 TTL 过期，然后转发到真正的处理队列
channel.queue_declare(
    queue="delay_30s",
    durable=True,
    arguments={
        "x-message-ttl": 30000,                    # 30 秒
        "x-dead-letter-exchange": "business_ex",   # 过期后转发目标
        "x-dead-letter-routing-key": "process",
    },
)

# 发布消息到延迟队列（不指定 exchange，直接用队列名）
channel.basic_publish(
    exchange="",
    routing_key="delay_30s",
    body=b"delayed task",
    properties=pika.BasicProperties(delivery_mode=2),
)

```

注意：此方案中延迟队列不应有消费者，否则消息会立即被消费而不会等待 TTL。

### 消息确认与重试策略

```python
import time

def on_message(channel, method, properties, body):
    try:
        process(body)
        channel.basic_ack(delivery_tag=method.delivery_tag)
    except TemporaryError:
        # 可恢复错误：重新入队，等待重试
        channel.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
    except PermanentError:
        # 不可恢复错误：拒绝消息，触发死信队列
        channel.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

```

重试次数控制：可在消息 headers 中记录重试次数，超过阈值后发送到死信队列，避免无限重试。

---

## 踩坑与注意事项

**忘记 ack 导致消息堆积**：使用手动 ack 时，若消费者处理消息后未调用 `basic_ack`，消息会一直处于 `Unacked` 状态，不会重新投递给其他消费者，直到连接断开后才重新入队。长期运行后 Unacked 消息数持续增长会导致内存告警。务必在 try/except 中确保 ack 或 nack 被调用。

**auto\_ack=True 丢消息**：设置 `auto_ack=True` 后，Broker 将消息投递给消费者的瞬间即视为已消费，无论消费者是否处理成功。消费者崩溃时当前批次消息全部丢失且无法恢复。生产环境应始终使用 `auto_ack=False` \+ 手动 ack。

**连接断开重连**：`pika.BlockingConnection` 不支持自动重连，网络抖动会抛出异常。生产环境的同步代码可捕获 `pika.exceptions.AMQPConnectionError` 后重建连接；异步场景应使用 `aio_pika.connect_robust()`。

**队列/Exchange 属性不可变**：已声明的 Queue 或 Exchange 的 `durable`、`type` 等属性无法修改。若需修改，必须先删除再重建，否则会抛出 `PRECONDITION_FAILED` 错误。

**消息体必须是 bytes**：`basic_publish` 的 `body` 参数只接受 `bytes`，传入 `str` 会报类型错误。序列化时应显式 `json.dumps(data).encode()` 或 `str.encode()`。

---

## 最佳实践

**Exchange + Routing Key 组合驱动路由，不要直接向 Queue 发送消息**：直接向 Queue 发送消息（使用默认 Exchange）绕过了 RabbitMQ 的路由机制，失去了 Topic/Fanout Exchange 的灵活性。生产代码应始终通过 Exchange 发送。

**声明队列和 Exchange 时使用 `durable=True`，消息使用 `delivery_mode=2`**：`durable=True` 保证 Broker 重启后队列/Exchange 不丢失；`delivery_mode=2`（持久化）保证消息在 Broker 崩溃时写入磁盘，两者缺一不可。

**消费者设置 `prefetch_count` 限制未确认消息数**：`basic_qos(prefetch_count=1)` 确保消费者在确认当前消息之前不接收新消息，防止一个消费者积压大量未处理消息，实现公平分发。

**使用死信队列（DLQ）处理失败消息，而非直接丢弃**：配置 `x-dead-letter-exchange` 后，超过重试次数或被拒绝（`nack + requeue=False`）的消息自动路由到死信队列，便于人工排查和重处理，避免消息静默丢失。

**异步场景使用 `aio_pika.connect_robust()`**：`aio_pika` 的 `connect_robust` 内置重连逻辑，网络抖动后自动恢复，不需要手动重建连接；同步场景用 `pika.BlockingConnection` 需自行处理重连。

---

## 常见陷阱

### 陷阱：忘记 ack 导致 Unacked 消息堆积

**现象：** RabbitMQ Management UI 中队列的 Unacked 消息数持续增长，消息积压，内存告警；重启消费者后 Unacked 消息重新入队。

**原因：** 使用 `auto_ack=False` 手动 ack 时，若消费者处理消息后未调用 `basic_ack`，消息一直处于 Unacked 状态，不会被重新分配，直到连接断开后才重入队。

**解决：** 在 try/finally 或 try/except 中确保 ack 或 nack 必定被调用。

```python
def on_message(channel, method, properties, body):
    try:
        process(body)
        channel.basic_ack(delivery_tag=method.delivery_tag)
    except Exception:
        channel.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

```

### 陷阱：修改已存在队列的属性导致 `PRECONDITION_FAILED`

**现象：** 消费者或生产者启动时抛出 `PRECONDITION_FAILED - inequivalent arg 'durable'`，连接被关闭。

**原因：** RabbitMQ 的 Queue 和 Exchange 一旦声明，`durable`、`type`、`x-*` 参数均不可修改。用不同参数重新声明同名队列会触发此错误。

**解决：** 先通过 Management UI 或 `rabbitmqadmin` 删除旧队列，再用新参数重建；或使用不同的队列名称。

### 陷阱：`auto_ack=True` 导致消息在消费者崩溃时丢失

**现象：** 消费者处理消息中途崩溃后，该批次消息丢失，无法从队列恢复。

**原因：** `auto_ack=True` 时，Broker 投递消息的瞬间即视为消费完成，无论消费者是否真正处理成功。

**解决：** 生产环境始终使用 `auto_ack=False` 加手动 ack，仅在消息处理成功后调用 `basic_ack`。

---

## 参见

- [Kafka完全指南](https://blog.vercanti.com/kafka-wan-quan-zhi-nan/)
- [Redis完全指南](https://blog.vercanti.com/redis-wan-quan-zhi-nan/)
- [FastAPI完全指南](https://blog.vercanti.com/fastapi-wan-quan-zhi-nan/)