RabbitMQ 完全指南
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 的区别
官方文档: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(同步)
安装
pip install pika
建立连接
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
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
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 |
发布消息
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 |
消费消息
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
# 确认单条消息
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)
公平分发
防止某个消费者堆积过多未处理消息:
# 每次只给消费者推送 1 条未 ack 的消息
channel.basic_qos(prefetch_count=1)
Python 客户端:aio-pika(异步)
安装
pip install aio-pika
异步连接
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。
异步发布消息
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",
)
异步消费消息
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 集成
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 重启后不丢失,需同时满足三个条件:
# 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 过期
- 队列达到最大长度
# 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 队列 + 死信队列"模拟:
# 延迟队列:消息在此队列中等待 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。
消息确认与重试策略
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 必定被调用。
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。