6.4 消息队列入门——RabbitMQ 与 Redis Stream
预计阅读时间:15 分钟
📖 目录
学习目标
- 理解消息队列的核心作用:异步解耦、削峰填谷、可靠投递
- 掌握 RabbitMQ 的 AMQP 模型:Exchange、Queue、Binding、Routing Key 四种类型
- 能用 Docker 部署 RabbitMQ,用 Python 编写生产-消费程序并配置死信队列
- 理解 Redis Stream 数据结构和消费者组机制
- 能根据场景正确选择 RabbitMQ / Redis Stream / Kafka
核心知识
| 主题 | 关键命令 / 接口 | 关联章节 |
|---|---|---|
| 消息队列模型 | 异步、缓冲、发布订阅 | — |
| RabbitMQ Exchange 类型 | Direct / Fanout / Topic / Headers | — |
| RabbitMQ 部署 | Docker Compose, rabbitmqadmin | 3.1:Docker 容器入门 Docker 入门 |
| 死信队列 | x-dead-letter-exchange, x-message-ttl | — |
| Redis Stream | XADD, XREAD, XREADGROUP, XACK | 3.10:Redis 缓存服务 Redis 缓存服务 |
| 消费者组 | XGROUP, XCLAIM, XPENDING | — |
知识关联
- 前置知识:3.1:Docker 容器入门 Docker 容器入门(Docker Compose 部署 RabbitMQ)、3.10:Redis 缓存服务 Redis 缓存服务(Redis Stream 前置知识)
- 后续影响:6.2:Kubernetes 入门 Kubernetes 入门(K8s 中部署消息队列集群)、4.9:CI/CD 持续部署 CI/CD 持续部署(异步解耦触发流水线)
- 配套技术:RabbitMQ + Redis Stream + Kafka 构成消息队列选型三角,AMQP 协议是 RabbitMQ 的通信基石
原理讲解
为什么消息队列要支持消息持久化
消息持久化是消息队列可靠性的基石。如果消息只存在于内存中,队列进程崩溃或服务器宕机后消息将永久丢失。持久化将消息写入磁盘(或追加日志),确保即使进程重启也能从磁盘恢复未处理的消息。RabbitMQ 通过将队列标记为 durable 并设置 delivery_mode=2 实现消息级持久化;Redis Stream 通过 AOF/RDB 持久化;Kafka 通过顺序写磁盘实现高吞吐持久化。持久化的代价是性能下降(磁盘 I/O 比内存慢 1000 倍),因此需要权衡:对于不允许丢消息的场景(订单、支付)必须开启持久化;对于允许少量丢失的场景(日志、监控)可以关闭。
RabbitMQ vs Kafka:架构哲学的根本差异
RabbitMQ 和 Kafka 代表了两种不同的消息模型。RabbitMQ 是"消息代理(Broker)"——生产者发消息到 Exchange,Exchange 路由到 Queue,消费者从 Queue 拉取并确认,消息被消费后删除。这种模型适合"任务分发"场景:每条消息只被一个消费者处理,处理完即删除。Kafka 是"分布式提交日志(Commit Log)"——生产者追加写入 Topic 的 Partition,消费者按 offset 顺序读取,消息在保留期内不删除。这种模型适合"事件流"场景:同一消息可被多个消费者独立处理(发布订阅),且支持历史回放。选型原则:需要灵活路由、任务队列、死信队列 → RabbitMQ;需要高吞吐、流式处理、事件溯源 → Kafka。
Redis Stream 为什么适合轻量级消息场景
Redis Stream 是 Redis 5.0 引入的追加日志数据结构,设计目标是在 Redis 内实现轻量级消息队列。相比 RabbitMQ,Stream 的优势是零额外运维(复用已有的 Redis 实例)、微秒级延迟、内存级性能。相比 Kafka,Stream 的劣势是吞吐量受限于单机内存、不支持多副本持久化。Stream 的消费者组机制提供了类似 Kafka 的负载均衡和消息确认能力(XREADGROUP + XACK),Pending 列表支持故障转移(XCLAIM 转移超时消息)。适用场景:已有 Redis 基础设施的轻量级消息需求、IoT 数据采集、实时事件处理。
3.1 消息队列通信模型
消息队列(Message Queue)是一种异步通信机制——生产者发送消息不等待处理结果,消费者异步处理。它解决的四个核心问题:
- 异步解耦:订单系统发消息后立即返回,邮件/短信/库存系统各自消费,互不阻塞
- 流量削峰:秒杀请求先涌入消息队列,后端按处理能力慢慢消费,不被打垮
- 发布订阅:一条消息被多个消费者独立处理(如一条日志同时进 ES 做搜索 + 进 Hadoop 做分析)
- 可靠投递:消息持久化 + 确认机制,系统崩溃后重启继续处理
3.2 RabbitMQ AMQP 模型
RabbitMQ 遵循 AMQP 0-9-1 协议,通信模型不是直接的"生产者→队列→消费者"。中间有一个关键角色——Exchange(交换机)。
# RabbitMQ 消息流
# 生产者 → Exchange → Binding → Queue → 消费者
# ↑
# Routing Key
Exchange 接收生产者的消息,根据 Routing Key 和 Binding 规则路由到目标 Queue。四种 Exchange 类型:
| Exchange 类型 | 路由逻辑 | 场景 |
|---|---|---|
| Direct | Routing Key 精确匹配 Binding Key | 单播,按级别路由日志 |
| Fanout | 广播到所有绑定队列 | 发布订阅,全局通知 |
| Topic | Routing Key 按模式匹配(通配符 * #) | 多条件路由,按类型+地域分发 |
| Headers | 按消息 Header 属性匹配 | 复杂的多条件路由(不常用) |
3.3 死信队列与 TTL
死信队列(Dead Letter Queue, DLQ)是处理消费失败消息的标准模式。当消息符合以下条件时,RabbitMQ 将其转发到指定的死信 Exchange:
- 消费者
basic.reject或basic.nack且requeue=false - 消息 TTL 过期(未在指定时间内消费)
- 队列达到最大长度,队列头的消息被丢弃
通过在声明队列时指定 x-dead-letter-exchange 参数来绑定 DLQ。死信 Exchange 将消息路由到死信队列,供后续人工介入或自动重试。
3.4 Redis Stream 数据结构
Redis 5.0+ 引入的 Stream 数据类型是一个追加日志(append-only log),每条消息有唯一 ID(格式:时间戳-序号)。相比 List 做队列,Stream 支持:
- 消费者组:多个消费者协同消费同一 Stream,每条消息只被组内一个消费者处理
- 消息确认:消费者处理完显式
XACK,未确认消息可重新分配 - 历史回溯:用
XRANGE 0从头重新消费
# Stream 结构示意
# ┌─────────────┬──────────────────────────────┐
# │ 消息 ID │ 键值对 (field-value pairs) │
# ├─────────────┼──────────────────────────────┤
# │ 1700000000000-0 │ sensor_id=1234 temperature=25.6│
# │ 1700000000001-0 │ sensor_id=5678 temperature=26.1│
# └─────────────┴──────────────────────────────┘
3.5 Redis Stream 消费者组模型
消费者组是 Stream 的核心特性。一个组内有多个消费者,每条消息被负载均衡到其中一个消费者。消费者组维护三个关键集合:
- 已交付未确认(Pending Entries List):消费者已读取但未
XACK的消息 - 游标(last_delivered_id):组内下一条待交付消息的 ID
- 消费者列表:组内所有消费者及其待处理消息数
# 消费者组工作流
# 生产者 → XADD → Stream → XREADGROUP → 消费者1 → XACK
# ↘ XREADGROUP → 消费者2 → XACK
# ↘ 消费者3 → XACK
示例代码
示例 1:Docker Compose 部署 RabbitMQ
cat <<'EOF' > docker-compose.rabbitmq.yml
services:
rabbitmq:
image: rabbitmq:4-management
container_name: rabbitmq
ports:
- "5672:5672" # AMQP 协议端口
- "15672:15672" # 管理界面
environment:
RABBITMQ_DEFAULT_USER: admin
RABBITMQ_DEFAULT_PASS: ${RABBITMQ_PASSWORD}
RABBITMQ_DEFAULT_VHOST: /
volumes:
- rabbitmq_data:/var/lib/rabbitmq
- rabbitmq_logs:/var/log/rabbitmq
restart: unless-stopped
healthcheck:
test: ["CMD", "rabbitmqctl", "status"]
interval: 10s
timeout: 5s
retries: 5
volumes:
rabbitmq_data:
rabbitmq_logs:
EOF
docker compose -f docker-compose.rabbitmq.yml up -d
# 访问管理界面: http://localhost:15672
# 用户: admin / 密码: $RABBITMQ_PASSWORD
# 创建用户与 vhost(生产环境推荐)
docker exec rabbitmq rabbitmqctl add_user app_user ${APP_PASSWORD}
docker exec rabbitmq rabbitmqctl set_user_tags app_user management
docker exec rabbitmq rabbitmqctl add_vhost /myapp
docker exec rabbitmq rabbitmqctl set_permissions -p /myapp app_user ".*" ".*" ".*"
docker exec rabbitmq rabbitmqctl list_vhosts
docker exec rabbitmq rabbitmqctl list_users
# 输出: Listing vhosts ...
# 输出: /myapp
示例 2:RabbitMQ Python 生产者
import pika
import json
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost',
credentials=pika.PlainCredentials('admin', 'password')))
channel = connection.channel()
channel.exchange_declare(exchange='task_exchange',
exchange_type='direct', durable=True)
channel.queue_declare(queue='task_queue', durable=True)
channel.queue_bind(exchange='task_exchange',
queue='task_queue', routing_key='task')
for i in range(10):
msg = json.dumps({"id": i, "task": f"process_data_{i}"})
channel.basic_publish(
exchange='task_exchange',
routing_key='task',
body=msg,
properties=pika.BasicProperties(
delivery_mode=2,
content_type='application/json'))
print(f"发送: {msg}")
connection.close()
# 输出: 发送: {"id": 0, "task": "process_data_0"}
# 输出: 发送: {"id": 1, "task": "process_data_1"}
示例 3:RabbitMQ Python 消费者(手动 ACK + QoS)
import pika
import time
def callback(ch, method, properties, body):
print(f"收到: {body}")
time.sleep(1)
print(" 处理完成")
ch.basic_ack(delivery_tag=method.delivery_tag)
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost',
credentials=pika.PlainCredentials('admin', 'password')))
channel = connection.channel()
channel.queue_declare(queue='task_queue', durable=True)
channel.basic_qos(prefetch_count=1)
channel.basic_consume(queue='task_queue',
on_message_callback=callback)
print("等待消息...")
channel.start_consuming()
# 输出: 等待消息...
# 输出: 收到: b'{"id": 0, "task": "process_data_0"}'
# 输出: 处理完成
示例 4:RabbitMQ 死信队列配置
import pika
# 建立连接(与示例 2/3 相同的连接参数)
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost',
credentials=pika.PlainCredentials('admin', 'password')))
channel = connection.channel()
# === 步骤 1:声明死信 Exchange 和死信队列 ===
# 死信 Exchange 接收被 reject/nack/TTL 过期的消息
channel.exchange_declare(exchange='dlx_exchange',
exchange_type='direct', durable=True)
channel.queue_declare(queue='dead_queue', durable=True)
channel.queue_bind(exchange='dlx_exchange',
queue='dead_queue', routing_key='dead') # routing_key 与下面的 x-dead-letter-routing-key 对应
# === 步骤 2:声明主队列,绑定死信 Exchange ===
# x-dead-letter-exchange: 消息被拒绝/过期后转发到哪个 Exchange
# x-dead-letter-routing-key: 转发时使用的 routing_key
# x-message-ttl: 队列级别 TTL(毫秒),超时未消费则进入死信
channel.queue_declare(
queue='task_queue',
durable=True,
arguments={
'x-dead-letter-exchange': 'dlx_exchange',
'x-dead-letter-routing-key': 'dead',
'x-message-ttl': 60000}) # 60 秒
# === 步骤 3:发送消息(两种过期方式) ===
# 方式一:消息级 TTL(优先级高于队列级 TTL)
channel.basic_publish(
exchange='task_exchange',
routing_key='task',
body='expirable message',
properties=pika.BasicProperties(
delivery_mode=2, # 持久化消息
expiration='60000')) # 消息级 TTL 60 秒
# 方式二:消费端 reject(手动触发死信)
# channel.basic_reject(delivery_tag=method.delivery_tag, requeue=False)
# === 步骤 4:消费死信队列 ===
def dead_letter_callback(ch, method, properties, body):
print(f"死信: {body}")
# 可在此处分析失败原因、告警、或写入数据库
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(queue='dead_queue',
on_message_callback=dead_letter_callback)
channel.start_consuming()
# 预期输出: 死信: b'expirable message' (发送后等待 60 秒出现)
示例 5:Redis Stream 基础命令
# 生产者:添加消息
XADD mystream * sensor_id 1234 temperature 25.6
# 输出: 1700000000000-0
XADD mystream * sensor_id 5678 temperature 26.1
# 输出: 1700000000001-0
# 读取所有消息
XRANGE mystream - +
# 输出: 1) 1) 1700000000000-0 2) sensor_id 1234 temperature 25.6
# 创建消费者组(从头开始消费)
XGROUP CREATE mystream mygroup 0
# 消费者读取新消息(COUNT 1:从组内最旧未投递消息开始)
XREADGROUP GROUP mygroup consumer1 COUNT 1 STREAMS mystream >
# 输出: 1) mystream 2) 1700000000000-0 sensor_id 1234 temperature 25.6
# 确认处理完成
XACK mystream mygroup 1700000000000-0
# 输出: (integer) 1
# 查看待处理(未 ACK)的消息——已全部确认
XPENDING mystream mygroup
# 输出: 1) (integer) 0
# 转移超时未确认的消息
XCLAIM mystream mygroup consumer2 60000 1700000000000-0
示例 6:Redis Stream 消费者组 Bash 脚本
#!/bin/bash
# stream_consumer.sh <consumer_id>
STREAM="mystream"
GROUP="mygroup"
ID=$1
redis-cli XGROUP CREATE $STREAM $GROUP $ MKSTREAM 2>/dev/null || true
echo "消费者 $ID 启动..."
while true; do
result=$(redis-cli XREADGROUP GROUP $GROUP "consumer_$ID" \
COUNT 1 BLOCK 0 STREAMS $STREAM ">")
if [ -n "$result" ]; then
msg_id=$(echo "$result" | grep -oP '(?<=\d-)\d+')
echo "处理消息: $result"
sleep 0.5
redis-cli XACK $STREAM $GROUP "$msg_id" > /dev/null
fi
done
# 使用方式(启动两个消费者协同消费):
# bash stream_consumer.sh 1 &
# bash stream_consumer.sh 2 &
$ 表示从 Stream 末尾开始(只收新消息)。> 表示只收未被当前消费组处理的消息。历史消息用 0 从头读取。常见错误
| 错误 | 原因 | 解决 |
|---|---|---|
| RabbitMQ 消费者收不到消息 | Exchange 未绑定 Queue,或 Routing Key 不匹配 | 检查管理界面 Binding 页,确认 Exchange 和 Queue 已正确绑定 |
PRECONDITION_FAILED 参数冲突 | 队列已存在但声明参数(durable/arguments)不一致 | 删除队列重建,或统一声明参数 |
| RabbitMQ 消息丢失 | 生产者未设置 delivery_mode=2 或队列非 durable | 队列声明 durable=True,消息设置 delivery_mode=2 |
| 死信队列无消息进入 | 死信 Exchange 类型与 routing_key 不匹配,或消费者未 reject | 确认 DLX 为 direct/fanout,测试 basic.reject(requeue=false) |
Redis Stream XREADGROUP 返回空 | 消费者组未正确创建,或用 $ 导致只读新消息 | 创建组时用 0 从头消费,或确认消息已写入 Stream |
| Stream 消费者重复消费 | 消费者处理完未调用 XACK | 确保处理逻辑完成后调用 XACK,必要时用 XCLAIM 转移超时消息 |
XGROUP CREATE 报 BUSYGROUP | 消费者组已存在,重复创建 | 加 MKSTREAM 或忽略错误 2>/dev/null |
最佳实践
| 实践 | 说明 |
|---|---|
| 始终通过 Exchange 发消息 | 不要直接发到 Queue,用 Exchange + Binding 解耦生产者和消费者 |
| 生产环境使用独立 vhost 隔离 | 每个应用/环境创建独立 vhost,避免队列名冲突和权限混乱 |
消费者设置 prefetch_count=1 | 公平分发,避免一个消费者堆积大量消息而其他消费者空闲 |
| 必须使用手动 ACK | 关闭 auto_ack,处理完成后再确认,防止消息丢失 |
| 为重要队列配置死信 Exchange | 消费失败的消息不丢失,进入 DLQ 供后续分析或重试 |
| Redis Stream 用消费者组而非单独 XREAD | 消费者组支持负载均衡、消息确认、故障转移,避免消息争抢 |
定期 XTRIM 限制 Stream 长度 | 防止 Stream 无限增长耗尽内存,用 MAXLEN ~ 10000 近似裁剪 |
| 监控 Pending 数量 | XPENDING 持续增长说明消费者出问题,需及时告警 |
练习题
- RabbitMQ 基础部署:用 Docker Compose 启动 RabbitMQ,创建 vhost
/testapp,创建 fanout Exchangetest.fanout和两个 Queueq1q2绑定到该 Exchange,用rabbitmqadmin publish发送一条消息,验证两个 Queue 都收到。 - Python 生产者消费者:编写 producer.py 向 direct Exchange 发送 20 条消息(routing_key=
task),consumer.py 手动 ACK 消费并打印耗时。验证 prefetch_count=1 的公平分发效果。 - 死信队列实验:配置队列 x-message-ttl=10000(10 秒),绑定死信 Exchange。发送一条消息,等待 10 秒后用
rabbitmqctl list_queues确认死信队列收到消息。 - Redis Stream 消费者组:启动两个消费者组脚本实例,向 Stream 添加 10 条消息,观察两个消费者交替处理。用
XPENDING和XINFO GROUPS查看状态。 - 选型分析:给出以下场景的 MQ 选型并说明理由:① 电商订单系统需要灵活路由;② IoT 传感器数据采集,已有 Redis 基础设施;③ 每天处理 10 亿条用户行为日志,需要流式回放给多个数据消费者。
点击查看答案
- 参考 docker-compose.yml 启动 RabbitMQ。用
rabbitmqadmin declare exchange name=test.fanout type=fanout和rabbitmqadmin declare queue name=q1创建资源,Binding 用rabbitmqadmin declare binding source=test.fanout destination=q1。rabbitmqadmin publish exchange=test.fanout payload="hello"后两个 queue 都应收到一条消息。 - Producer 用
channel.basic_publish(exchange='', routing_key='task', body=msg)。Consumer 用channel.basic_qos(prefetch_count=1)实现公平分发,手动 ACK 在callback函数末尾调用ch.basic_ack(delivery_tag=method.delivery_tag)。观察两个 consumer 实例交替收到消息。 - 声明 queue 时加参数
arguments={"x-message-ttl": 10000, "x-dead-letter-exchange": "dlx"}。10 秒后原队列消息消失,死信队列出现消息。 - 两个 consumer 用相同消费者组名
XGROUP。XREADGROUP 读取时每条消息只投递给组内一个成员。XINFO GROUPS 显示待处理和挂起数量。 - ① 选 RabbitMQ(灵活路由+死信+确认机制);② 选 Redis Stream(已有 Redis 零额外运维);③ 选 Kafka(高吞吐+日志持久化+多消费者 replay)。
学习检查点
学完本章后,请检验自己是否掌握以下内容:
| 检查项 | 自测问题 | 验证方法 |
|---|---|---|
| 概念理解 | 能用自己的话解释消息队列的异步解耦、削峰填谷原理 | 尝试向他人讲解 |
| 命令操作 | 能不查文档完成 RabbitMQ/Redis Stream 的部署和消息收发 | 在终端实际执行 |
| 原理掌握 | 能说出消息队列的持久化机制和消费者组工作原理 | 画出流程图 |
| 故障排查 | 能独立排查消息丢失或消费者堆积的问题 | 模拟故障并修复 |
| 最佳实践 | 能说明为什么需要配置消息队列的死信队列和重试机制 | 对比不同方案 |
本章总结
速查表
| 概念 | 关键命令 / 参数 | 一句话总结 |
|---|---|---|
| 消息队列 | 异步通信 | 解耦生产者和消费者,削峰填谷 |
| RabbitMQ Exchange | direct / fanout / topic / headers | 消息路由器,按类型决定投递规则 |
| Binding | queue_bind, routing_key | Exchange 到 Queue 的连接,附带路由键 |
| 死信队列 (DLQ) | x-dead-letter-exchange | 消费失败或过期的消息转入专用队列 |
| 手动 ACK | basic_ack, XACK | 消费完成才确认,崩溃后可重新投递 |
| Redis Stream | XADD / XREAD / XREADGROUP | 追加日志结构,天然支持消费者组 |
| 消费者组 | XGROUP CREATE, XACK | 负载均衡 + 消息确认 + 故障转移 |
| Pending 列表 | XPENDING, XCLAIM | 未确认消息记录,支持超时重新分配 |
| 选型决策 | RabbitMQ vs Redis vs Kafka | 灵活路由→RabbitMQ;已有Redis→Stream;海量事件→Kafka |
RabbitMQ vs Redis Stream vs Kafka 选型对比表
| 维度 | RabbitMQ | Redis Stream | Kafka |
|---|---|---|---|
| 协议 | AMQP 0-9-1 | Redis 原生命令 | Kafka 协议(自定义 TCP) |
| 消息模型 | Exchange + Queue + Binding | 追加日志 + 消费者组 | Topic + Partition + Consumer Group |
| 消息持久化 | 队列级持久化(磁盘) | AOF/RDB 持久化 | 磁盘顺序写(高吞吐) |
| 消费确认 | basic_ack / basic.nack | XACK | 自动提交或手动提交 offset |
| 消息回溯 | 不支持(消费即删除) | 支持(XRANGE 从头读) | 支持(保留期内任意 offset 回放) |
| 吞吐量 | 万级/秒 | 十万级/秒 | 百万级/秒 |
| 延迟 | 微秒级 | 微秒级 | 毫秒级 |
| 运维复杂度 | 中等(Erlang OTP) | 低(复用 Redis) | 高(ZooKeeper/KRaft + 多 Broker) |
| 适用场景 | 复杂路由、企业级消息、任务队列 | 已有 Redis、轻量级流处理、IoT | 海量事件流、日志聚合、流式回放 |