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-ttlx-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_robustconnect 的区别: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_rejectbasic_nackrequeue=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 的 durabletype 等属性无法修改。若需修改,必须先删除再重建,否则会抛出 PRECONDITION_FAILED 错误。

消息体必须是 bytesbasic_publishbody 参数只接受 bytes,传入 str 会报类型错误。序列化时应显式 json.dumps(data).encode()str.encode()


最佳实践

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

声明队列和 Exchange 时使用 durable=True,消息使用 delivery_mode=2durable=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_pikaconnect_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 一旦声明,durabletypex-* 参数均不可修改。用不同参数重新声明同名队列会触发此错误。

解决: 先通过 Management UI 或 rabbitmqadmin 删除旧队列,再用新参数重建;或使用不同的队列名称。

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

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

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

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


参见

阅读更多

Web 安全基础

1. HTML 转义(服务端渲染必须): 2. CSP(Content Security Policy): 3. HttpOnly Cookie:防止 JS 读取会话 Cookie: 4. 前端框架防护: 攻击者在第三方网站构造一个表单,诱导已登录用户提交,浏览器会自动携带目标站的 Cookie。 触发条件: 1. 用户已登录目标网站(Cookie 有效) 2. 目标 API 仅凭 Cookie 识别用户身份 3. 请求来源未验证 1. CSRF Token(推荐): 2. SameSite Cookie: 3. 验证 Origin/Referer 头:

By yellowdog

HTTP 协议深度指南

HTTP(HyperText Transfer Protocol)是 Web 的基础传输协议,基于 TCP/IP,采用请求/响应模型。 相关文档:Web安全基础(/web-an-quan-ji-chu/) FastAPI完全指南(/fastapi-wan-quan-zhi-nan/) Nginx完全指南(/nginx-wan-quan-zhi-nan/) 幂等性:多次执行相同请求,服务器状态结果相同。PUT /users/1 多次执行结果一致;POST /users 每次创建新资源,非幂等。 浏览器直接从本地缓存读取,不向服务器发送请求。 缓存命中时,状

By yellowdog

系统设计基础

SLA 对照表: 选择建议:无状态服务(Web 层、API 层)优先水平扩展;数据库初期垂直扩展,达到瓶颈后考虑分库分表或读写分离。 缓存穿透(查询不存在的 key,每次都打到 DB): 缓存击穿(热点 key 过期,瞬间大量请求打到 DB): 缓存雪崩(大量 key 同时过期,或缓存服务宕机): 令牌桶 Python 实现: Redis 实现分布式限流(滑动窗口): URL 命名规则: Cursor 分页响应格式: 雪花算法结构(64 bit): 定义:分布式系统不能同时满足以下三个特性: 在分布式环境中 P 是必须保证的,所以实际是 CP vs AP

By yellowdog

算法思路与模板

二分查找要求序列有序,每次将搜索范围缩减一半,时间复杂度 O(log n)。 两个指针从两端向中间收缩,常用于有序数组。 滑动窗口维护一个满足条件的区间 left, right,right 不断向右扩张,条件不满足时收缩 left。 滑动窗口通用框架: 1. 确定"子问题":原问题可以分解为哪些规模更小的同类问题 2. 定义 dpi 或 dpij 的含义,要足够清晰 3. 推导状态转移方程 4. 确定初始状态(边界条件) 5. 确定计算顺序(确保依赖的子问题先计算) 每件物品最多选一次。dpj = 容量为 j 时的最大价值,逆序遍历容量防止重复选取。 每

By yellowdog