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, rabbitmqadmin3.1:Docker 容器入门 Docker 入门
死信队列x-dead-letter-exchange, x-message-ttl
Redis StreamXADD, XREAD, XREADGROUP, XACK3.10:Redis 缓存服务 Redis 缓存服务
消费者组XGROUP, XCLAIM, XPENDING

知识关联

原理讲解

为什么消息队列要支持消息持久化

消息持久化是消息队列可靠性的基石。如果消息只存在于内存中,队列进程崩溃或服务器宕机后消息将永久丢失。持久化将消息写入磁盘(或追加日志),确保即使进程重启也能从磁盘恢复未处理的消息。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 KeyBinding 规则路由到目标 Queue。四种 Exchange 类型:

Exchange 类型路由逻辑场景
DirectRouting Key 精确匹配 Binding Key单播,按级别路由日志
Fanout广播到所有绑定队列发布订阅,全局通知
TopicRouting Key 按模式匹配(通配符 * #多条件路由,按类型+地域分发
Headers按消息 Header 属性匹配复杂的多条件路由(不常用)
⚠️ 队列 vs Exchange 初学者常犯的错误:直接把消息发到 Queue 上。正确的做法是:发到 Exchange,由 Binding 路由到 Queue。Exchange 是路由器,Queue 是信箱。

3.3 死信队列与 TTL

死信队列(Dead Letter Queue, DLQ)是处理消费失败消息的标准模式。当消息符合以下条件时,RabbitMQ 将其转发到指定的死信 Exchange:

  • 消费者 basic.rejectbasic.nackrequeue=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 &
💡 $ vs > $ 表示从 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 CREATEBUSYGROUP消费者组已存在,重复创建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 持续增长说明消费者出问题,需及时告警

练习题

  1. RabbitMQ 基础部署:用 Docker Compose 启动 RabbitMQ,创建 vhost /testapp,创建 fanout Exchange test.fanout 和两个 Queue q1 q2 绑定到该 Exchange,用 rabbitmqadmin publish 发送一条消息,验证两个 Queue 都收到。
  2. Python 生产者消费者:编写 producer.py 向 direct Exchange 发送 20 条消息(routing_key=task),consumer.py 手动 ACK 消费并打印耗时。验证 prefetch_count=1 的公平分发效果。
  3. 死信队列实验:配置队列 x-message-ttl=10000(10 秒),绑定死信 Exchange。发送一条消息,等待 10 秒后用 rabbitmqctl list_queues 确认死信队列收到消息。
  4. Redis Stream 消费者组:启动两个消费者组脚本实例,向 Stream 添加 10 条消息,观察两个消费者交替处理。用 XPENDINGXINFO GROUPS 查看状态。
  5. 选型分析:给出以下场景的 MQ 选型并说明理由:① 电商订单系统需要灵活路由;② IoT 传感器数据采集,已有 Redis 基础设施;③ 每天处理 10 亿条用户行为日志,需要流式回放给多个数据消费者。
点击查看答案
  1. 参考 docker-compose.yml 启动 RabbitMQ。用 rabbitmqadmin declare exchange name=test.fanout type=fanoutrabbitmqadmin declare queue name=q1 创建资源,Binding 用 rabbitmqadmin declare binding source=test.fanout destination=q1rabbitmqadmin publish exchange=test.fanout payload="hello" 后两个 queue 都应收到一条消息。
  2. 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 实例交替收到消息。
  3. 声明 queue 时加参数 arguments={"x-message-ttl": 10000, "x-dead-letter-exchange": "dlx"}。10 秒后原队列消息消失,死信队列出现消息。
  4. 两个 consumer 用相同消费者组名 XGROUP。XREADGROUP 读取时每条消息只投递给组内一个成员。XINFO GROUPS 显示待处理和挂起数量。
  5. ① 选 RabbitMQ(灵活路由+死信+确认机制);② 选 Redis Stream(已有 Redis 零额外运维);③ 选 Kafka(高吞吐+日志持久化+多消费者 replay)。

学习检查点

学完本章后,请检验自己是否掌握以下内容:

检查项自测问题验证方法
概念理解能用自己的话解释消息队列的异步解耦、削峰填谷原理尝试向他人讲解
命令操作能不查文档完成 RabbitMQ/Redis Stream 的部署和消息收发在终端实际执行
原理掌握能说出消息队列的持久化机制和消费者组工作原理画出流程图
故障排查能独立排查消息丢失或消费者堆积的问题模拟故障并修复
最佳实践能说明为什么需要配置消息队列的死信队列和重试机制对比不同方案

本章总结

速查表

概念关键命令 / 参数一句话总结
消息队列异步通信解耦生产者和消费者,削峰填谷
RabbitMQ Exchangedirect / fanout / topic / headers消息路由器,按类型决定投递规则
Bindingqueue_bind, routing_keyExchange 到 Queue 的连接,附带路由键
死信队列 (DLQ)x-dead-letter-exchange消费失败或过期的消息转入专用队列
手动 ACKbasic_ack, XACK消费完成才确认,崩溃后可重新投递
Redis StreamXADD / XREAD / XREADGROUP追加日志结构,天然支持消费者组
消费者组XGROUP CREATE, XACK负载均衡 + 消息确认 + 故障转移
Pending 列表XPENDING, XCLAIM未确认消息记录,支持超时重新分配
选型决策RabbitMQ vs Redis vs Kafka灵活路由→RabbitMQ;已有Redis→Stream;海量事件→Kafka

RabbitMQ vs Redis Stream vs Kafka 选型对比表

维度RabbitMQRedis StreamKafka
协议AMQP 0-9-1Redis 原生命令Kafka 协议(自定义 TCP)
消息模型Exchange + Queue + Binding追加日志 + 消费者组Topic + Partition + Consumer Group
消息持久化队列级持久化(磁盘)AOF/RDB 持久化磁盘顺序写(高吞吐)
消费确认basic_ack / basic.nackXACK自动提交或手动提交 offset
消息回溯不支持(消费即删除)支持(XRANGE 从头读)支持(保留期内任意 offset 回放)
吞吐量万级/秒十万级/秒百万级/秒
延迟微秒级微秒级毫秒级
运维复杂度中等(Erlang OTP)低(复用 Redis)高(ZooKeeper/KRaft + 多 Broker)
适用场景复杂路由、企业级消息、任务队列已有 Redis、轻量级流处理、IoT海量事件流、日志聚合、流式回放

延伸阅读

常见问题

RabbitMQ 的死信队列是做什么的?
死信队列(Dead Letter Queue)用于存储无法被正常消费的消息。消息成为死信的三种情况:① 消费者使用 basic.reject/basic.nack 且 requeue=false;② 消息 TTL 过期;③ 队列达到最大长度。死信可被重新路由到指定 Exchange 做二次分析或人工处理。生产环境必须配置死信队列,避免消息丢失。
Redis Stream 和 Kafka 的核心区别?
Redis Stream 是内存数据结构,数据存在 Redis 中,吞吐量受限于单机内存。Kafka 是分布式流平台,数据持久化到磁盘,支持分区扩展、多消费者组、数据长时间保留。Redis Stream 适合轻量级消息(通知、日志汇总),Kafka 适合企业级事件流(日志聚合、CDC 数据同步)。
RabbitMQ 消息积压怎么处理?
先排查消费者是否 hang 住或处理太慢。临时方案:① 增加消费者数量(扩容);② 临时调高 prefetch_count 增加消费者缓冲区;③ 将消息转存到文件或另一队列做延迟处理。长期方案:优化消费者处理逻辑、增加 Queue 分区数。如果消费者一直赶不上,说明架构需要调整(如增加 Exchange)。
↑ 回到顶部