Redis 不只是缓存数据库,它在消息队列领域也占据重要位置。从早期的 Pub/Sub 发布订阅,到 Redis 5.0 引入的 Streams 流数据结构,Redis 提供了两种截然不同但各有适用场景的消息传递模型。本文将深入对比这两种方案的原理、命令、风险与实战用法,并给出与 Kafka、RabbitMQ 的专业选型对比。
一、Pub/Sub:轻量广播模型
Redis Pub/Sub(Publish/Subscribe)是最早支持的消息机制,采用经典的发布订阅模式。它的设计理念是"fire and forget":消息发布后立即推送给所有在线订阅者,不做任何持久化。
1.1 核心命令
# 客户端 A 订阅频道
SUBSCRIBE orders
# 客户端 B 订阅匹配模式(通配符)
PSUBSCRIBE orders.*
# 发布者推送消息
PUBLISH orders '{"order_id":"20250813001","status":"paid","amount":299.00}'
# 查看当前被订阅的频道(服务端)
PUBSUB CHANNELS
PUBSUB NUMSUB orders
SUBSCRIBE 订阅指定频道,PSUBSCRIBE 通过 glob 模式批量订阅(如 orders.* 匹配 orders.paid、orders.shipped 等),PUBLISH 向频道广播消息。注意:所有命令在订阅状态下会阻塞 Redis 连接,实际应用中通常使用独立连接处理订阅,另一个连接处理其他 Redis 操作。
1.2 消息传递机制
Pub/Sub 的关键特征是推(Push)模型:
- 发布者调用
PUBLISH channel message - Redis Server 遍历该频道的所有订阅者列表
- 将消息立即推送到每个订阅者的客户端缓冲区
- 如果订阅者网络延迟或处理缓慢,消息堆积在服务器输出缓冲区
Publisher Redis Server Subscriber A
| | |
|-- PUBLISH orders ---> | |
| |-- 推送给所有 subscriber --->|
| | | 收到消息
| |-- 推送给所有 subscriber --->| Subscriber B
1.3 Pub/Sub 的致命缺陷:消息丢失
Pub/Sub 有三个导致消息丢失的典型场景,这也是它在生产环境被 Streams 逐步替代的核心原因。
场景一:订阅者离线
# subscriber 断开连接期间发布的消息全部丢失
# Redis 不保存消息历史,新上线的 subscriber 无法收到离线期间的消息
场景二:输出缓冲区溢出
# Redis 配置中的 client-output-buffer-limit pubsub
# 默认:client-output-buffer-limit pubsub 32mb 8mb 60
# 当 subscriber 消费慢于生产速度,超过 32MB 硬限制时
# Redis 会强制断开该 subscriber 连接,期间消息全部丢失
场景三:网络分区
Publisher -> [网络分区] -> Redis <-> Subscriber
# 分区期间发布的消息,subscriber 永远收不到
1.4 适用场景
尽管有丢失风险,Pub/Sub 在以下场景仍非常高效:
| 场景 | 说明 |
|---|---|
| 实时通知 | 在线用户的消息提醒、弹幕推送 |
| 配置热更新 | 配置中心广播配置变更信号 |
| 缓存失效广播 | 分布式缓存一致性通知 |
| 日志聚合 | 实时日志流(允许丢失少量日志) |
| 心跳检测 | 服务健康状态广播 |
Pub/Sub 的优势在于极低延迟(亚毫秒级)和极简实现(无磁盘 I/O),在允许偶发丢失且 subscriber 必须实时在线的场景下,它是最佳选择。
# 典型缓存失效广播用法
PUBLISH cache:invalidate "user:profile:10086"
# 所有缓存节点同时收到,立即删除本地缓存
二、Streams:持久化日志型消息队列
Redis 5.0 引入的 Streams 是 Redis 最重要的数据结构之一,设计灵感直接来自 Apache Kafka。它将消息以追加写(Append-Only)方式存储在内存中(可配置持久化到 RDB/AOF),每个消息拥有全局唯一、单调递增的 ID。
2.1 Streams 核心数据结构
Stream Key: orders:events
+------------------------------------------+
| ID (毫秒时间戳-序列号) | Field-Value |
+------------------------------------------+
| 1723512000000-0 | order_id 001 |
| | status paid |
| | amount 299.00 |
+------------------------------------------+
| 1723512000523-0 | order_id 002 |
| | status shipped|
+------------------------------------------+
| 1723512000891-0 | order_id 003 |
| | status refund |
+------------------------------------------+
Stream ID 格式为 millisecondsTime-sequenceNumber,由 Redis 自动分配,保证全局有序且唯一。也可以用 * 让 Redis 自动生成,或自定义 ID(必须比已有 ID 大)。
2.2 生产者:XADD
# 自动分配 ID(推荐)
XADD orders:events * order_id 001 status paid amount 299.00
# 返回: 1723512000000-0
# 自定义 ID(必须递增)
XADD orders:events 1723512000100-0 order_id 002 status shipped
# 限制 Stream 长度(近似裁剪,类似 Kafka retention)
XADD orders:events MAXLEN ~ 10000 * order_id 003 status refund
# 严格限制长度(精确裁剪,性能略低)
XADD orders:events MAXLEN 10000 * order_id 004 status completed
MAXLEN ~ 使用近似裁剪,通过 radix tree 快速定位删除节点,性能损耗极小(推荐生产使用)。~ 表示允许少量超出目标长度。
2.3 消费者:XREAD
# 阻塞读取新消息($ 表示只接收执行后产生的新消息)
XREAD BLOCK 5000 STREAMS orders:events $
# 阻塞 5000ms,等待新消息到达
# 读取历史消息(从指定 ID 之后开始)
XREAD COUNT 10 STREAMS orders:events 1723512000000-0
# 同时监听多个 Stream
XREAD BLOCK 0 STREAMS orders:events payments:events $ $
# BLOCK 0 表示永久阻塞
# 作为消费者组中的独立消费者读取(见第三节)
XREADGROUP GROUP order_consumers consumer-1 COUNT 5 BLOCK 3000 \
STREAMS orders:events >
XREAD 是非破坏性地读取消息,消息仍保留在 Stream 中,可被多个消费者重复读取。这是 Streams 相对于传统 list 结构(BLPOP/BRPOP)的重大优势。
2.4 消息查询:XRANGE / XREVRANGE / XLEN
# 按 ID 范围查询
XRANGE orders:events 1723512000000-0 1723512000891-0
# 分页查询(- 表示最小 ID,+ 表示最大 ID)
XRANGE orders:events - + COUNT 10
# 反向查询(从最新开始)
XREVRANGE orders:events + - COUNT 5
# 获取 Stream 长度
XLEN orders:events
# 删除特定消息
XDEL orders:events 1723512000000-0
# 查看 Stream 信息
XINFO STREAM orders:events
XINFO GROUPS orders:events
XINFO CONSUMERS orders:events order_consumers
2.5 消息裁剪与清理
# 手动裁剪到指定长度
XTRIM orders:events MAXLEN ~ 5000
# 删除整个 Stream
DEL orders:events
# 基于 ID 范围删除(比如删除 7 天前的消息)
# 先找到 7 天前的毫秒时间戳,然后删除区间
XRANGE orders:events - 1722907200000-0 COUNT 1000
# 遍历结果并用 XDEL 逐个删除
生产环境通常配置 MAXLEN ~ 自动裁剪,或配合定时任务基于时间范围清理,模拟 Kafka 的消息保留策略。
三、Consumer Group:水平扩展与 ACK 确认
Consumer Group(消费者组)是 Streams 实现可靠消费的核心机制,功能直接对标 Kafka 的 Consumer Group。
3.1 核心概念
Stream: orders:events
Consumer Group: group1
+--------------------+
| Consumer A |-- 消费 partition 1(消息 ID: 1-5, 11-15...)
| Consumer B |-- 消费 partition 2(消息 ID: 6-10, 16-20...)
| Consumer C |-- 消费 partition 3(消息 ID: 21-25...)
+--------------------+
Pending Entries List (PEL):
+----------------+----------------+---------------+
| 消息 ID | 消费者 | 投递时间 |
+----------------+----------------+---------------+
| 1723512000000-0| consumer-a | 1723512005000 |
| 1723512000523-0| consumer-b | 1723512005500 |
+----------------+----------------+---------------+
关键设计:
- 消息不删除:每个消息被组内一个消费者接收,但消息本身仍在 Stream 中
- ACK 确认:消费者处理完成后必须发送 XACK,否则消息留在 PEL(待处理列表)
- 故障转移:消费者宕机后,其他消费者可以用
XCLAIM接管其待处理消息 - 游标管理:Redis 为每个消费者组维护最后已交付的消息 ID,新消费者从该位置继续
3.2 消费者组管理命令
# 1. 创建消费者组(从 Stream 起始位置开始)
XGROUP CREATE orders:events order_consumers 0
# 0 表示从头消费,$ 表示只消费创建后的新消息
# 2. 读取消息(> 表示读取未分配给任何消费者的新消息)
XREADGROUP GROUP order_consumers consumer-1 COUNT 5 BLOCK 3000 \
STREAMS orders:events >
# 3. 确认消息已处理
XACK orders:events order_consumers 1723512000000-0 1723512000523-0
# 4. 查看待处理消息(PEL)
XPENDING orders:events order_consumers - + 10
# 5. 查看某个消费者的待处理消息
XPENDING orders:events order_consumers - + 10 consumer-1
# 6. 转移未确认消息(consumer-1 宕机后,consumer-2 接管)
XCLAIM orders:events order_consumers consumer-2 60000 \
1723512000000-0
# 60000 是 idle time,只转移空闲超过 60 秒的消息
3.3 消费者组消费流程
# 完整消费循环示例
# Step 1: 创建组(仅需执行一次)
XGROUP CREATE orders:events order_group $ MKSTREAM
# MKSTREAM: 如果 Stream 不存在则自动创建
# Step 2: 消费者读取消息
XREADGROUP GROUP order_group worker-1 COUNT 1 BLOCK 5000 \
STREAMS orders:events >
# 返回:
# 1) 1) "orders:events"
# 2) 1) 1) "1723512000000-0"
# 2) 1) "order_id"
# 2) "001"
# 3) "status"
# 4) "paid"
# Step 3: 业务处理...
# 处理完成后确认
XACK orders:events order_group 1723512000000-0
# Step 4: 如果处理失败不发送 XACK,消息留在 PEL
# 通过监控 PEL 发现未处理消息并重试
3.4 待处理消息监控与故障恢复
# 查看 PEL 概览
XPENDING orders:events order_group
# 返回: 1) (integer) 3 -- 待处理消息数
# 2) "1723512000000-0" -- 最小 ID
# 3) "1723512000891-0" -- 最大 ID
# 4) 1) 1) "worker-1" -- 消费者
# 2) "2" -- 该消费者持有的待处理数
# 查看具体待处理消息详情
XPENDING orders:events order_group - + 10
# 返回每个消息的 ID、消费者、空闲时间、投递次数
# 自动转移长时间未处理的消息(idel 超过 60 秒)
XPENDING orders:events order_group - + 10
# 拿到待处理 ID 列表后
XCLAIM orders:events order_group worker-backup 60000 \
1723512000000-0 1723512000523-0
# 更现代的方式:XAUTOCLAIM(Redis 6.2+)
XAUTOCLAIM orders:events order_group worker-backup 60000 - \
COUNT 100
# 自动转移所有 idle 超过 60 秒的待处理消息
3.5 Streams vs List 作为队列
Redis 早期使用 List(LPUSH/BRPOP)实现队列,Streams 在功能上全面超越:
| 特性 | List (LPUSH/BRPOP) | Streams |
|---|---|---|
| 持久化 | 无专用机制 | 内建日志结构 |
| 消息 ID | 无 | 全局有序时间戳 ID |
| 多消费者 | 互斥竞争 | 组内负载均衡 |
| 消息确认 | 不支持 | XACK 确认 |
| 消息追踪 | 不支持 | PEL 待处理列表 |
| 历史查询 | 不支持 | XRANGE 范围查询 |
| 阻塞读取 | BRPOP | XREAD BLOCK |
| 内存效率 | 高 | 中(附加元数据开销) |
结论:所有新开发的消息队列场景都应使用 Streams,List 队列已进入维护模式。
四、Kafka / RabbitMQ / Redis Streams 选型对比
| 维度 | Redis Streams | Apache Kafka | RabbitMQ |
|---|---|---|---|
| 数据持久化 | 内存为主,RDB/AOF 落盘 | 磁盘持久化,页缓存优化 | 内存/磁盘可选 |
| 消息保留 | MAXLEN/~ 手动或自动裁剪 | 时间/大小策略,长期保留 | TTL/队列长度限制 |
| 吞吐量 | 中高(~100K msg/s 单机) | 极高(百万级/秒集群) | 中高(~50K msg/s) |
| 延迟 | 极低(亚毫秒) | 低(毫秒级) | 低(毫秒级) |
| 消费者模型 | Push + Pull | Pull | Push + Pull |
| Consumer Group | 支持 | 原生支持 | 需插件或手动实现 |
| 消息回溯 | 支持(XRANGE) | 原生支持(offset 回溯) | 有限支持 |
| 消息顺序 | Stream 内严格有序 | Partition 内有序 | Queue 内有序 |
| 水平扩展 | Redis Cluster 分片 | 原生分区扩展 | Cluster/ Federation |
| 运维复杂度 | 低 | 高 | 中 |
| 适用场景 | 中等规模实时流 | 大规模日志/事件流 | 企业级消息路由 |
4.1 选型建议
选择 Redis Streams 当:
- 消息量 < 100万/天,或 消费者延迟要求 < 10ms
- 不想引入额外基础设施(已部署 Redis)
- 消息保留期较短(小时到几天级别)
- 需要简单查询历史消息
选择 Kafka 当:
- 消息量 > 100万/天,或需要长期保留(周/月级别)
- 需要多消费者组独立消费同一数据
- 需要流处理(Kafka Streams / Flink)
- 团队有 Kafka 运维能力
选择 RabbitMQ 当:
- 需要复杂路由(Topic、Headers、Direct)
- 需要事务消息、死信队列等企业特性
- 需要 AMQP 协议兼容性
4.2 Redis Streams 的典型反模式
# 反模式 1:用 Streams 替代 Kafka 处理海量日志
# 问题:Redis 内存成本远高于磁盘,MAXLEN 频繁裁剪影响性能
# 替代:Kafka + 冷存(S3/对象存储)
# 反模式 2:一个 Stream 配过多 Consumer Group
# 问题:每个 Group 独立维护 PEL 和游标,内存开销线性增长
# 替代:控制 Group 数量(通常 < 10),或拆分 Stream
# 反模式 3:从不发送 XACK
# 问题:PEL 无限增长,消费者重启后重复消费大量旧消息
# 修正:处理成功立即 XACK,失败时记录日志后也要 XACK(避免阻塞)
# 反模式 4:BLOCK 0 永久阻塞不设超时
# 问题:连接断开检测延迟,影响故障恢复
# 修正:BLOCK 3000-5000,应用层循环重连
五、实战:订单事件流处理(Streams + Go)
本节实现一个完整的订单事件流系统,包含生产者、消费者组、ACK 确认和监控。
5.1 项目结构
order-stream/
├── producer.go # 订单事件生产者
├── consumer.go # 消费者组 Worker
├── monitor.go # PEL 监控与死信处理
├── docker-compose.yml # Redis 环境
└── go.mod
5.2 go.mod
module order-stream
go 1.22
require github.com/redis/go-redis/v9 v9.5.1
require (
github.com/cespare/xxhash/v2 v2.2.0 // indirect
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect
)
5.3 生产者:订单事件入队
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
"github.com/redis/go-redis/v9"
)
// OrderEvent 订单事件
type OrderEvent struct {
OrderID string `json:"order_id"`
UserID int64 `json:"user_id"`
Status string `json:"status"` // created, paid, shipped, delivered, refunded
Amount float64 `json:"amount"`
Timestamp time.Time `json:"timestamp"`
}
var ctx = context.Background()
// Producer 订单事件生产者
type Producer struct {
client *redis.Client
stream string
}
// NewProducer 创建生产者
func NewProducer(addr string) *Producer {
rdb := redis.NewClient(&redis.Options{
Addr: addr,
Password: "",
DB: 0,
PoolSize: 10,
})
return &Producer{
client: rdb,
stream: "orders:events",
}
}
// PublishEvent 发布订单事件到 Stream
func (p *Producer) PublishEvent(event OrderEvent) (string, error) {
data, err := json.Marshal(event)
if err != nil {
return "", fmt.Errorf("marshal event failed: %w", err)
}
// XADD 自动分配 ID,限制 Stream 最大长度约 10000 条
id, err := p.client.XAdd(ctx, &redis.XAddArgs{
Stream: p.stream,
MaxLen: 10000,
Approx: true,
Values: map[string]interface{}{
"data": string(data),
},
}).Result()
if err != nil {
return "", fmt.Errorf("xadd failed: %w", err)
}
log.Printf("[Producer] Event published, id=%s, order=%s, status=%s",
id, event.OrderID, event.Status)
return id, nil
}
// PublishBatch 批量发布(减少 RTT)
func (p *Producer) PublishBatch(events []OrderEvent) error {
pipe := p.client.Pipeline()
for _, event := range events {
data, _ := json.Marshal(event)
pipe.XAdd(ctx, &redis.XAddArgs{
Stream: p.stream,
MaxLen: 10000,
Approx: true,
Values: map[string]interface{}{
"data": string(data),
},
})
}
_, err := pipe.Exec(ctx)
return err
}
// Close 关闭连接
func (p *Producer) Close() error {
return p.client.Close()
}
// 演示:模拟订单事件流
func main() {
producer := NewProducer("localhost:6379")
defer producer.Close()
statuses := []string{"created", "paid", "shipped", "delivered"}
for i := 0; i < 100; i++ {
event := OrderEvent{
OrderID: fmt.Sprintf("ORD-%06d", i+1),
UserID: int64(10000 + i%100),
Status: statuses[i%len(statuses)],
Amount: 99.00 + float64(i)*10,
Timestamp: time.Now(),
}
if _, err := producer.PublishEvent(event); err != nil {
log.Printf("publish failed: %v", err)
}
time.Sleep(100 * time.Millisecond)
}
log.Println("All events published")
}
5.4 消费者:消费者组 Worker
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"os"
"os/signal"
"syscall"
"time"
"github.com/redis/go-redis/v9"
)
// Consumer 消费者组 Worker
type Consumer struct {
client *redis.Client
stream string
group string
consumerName string
batchSize int64
blockMS int64
}
// NewConsumer 创建消费者
func NewConsumer(addr, group, consumerName string) *Consumer {
rdb := redis.NewClient(&redis.Options{
Addr: addr,
PoolSize: 10,
})
return &Consumer{
client: rdb,
stream: "orders:events",
group: group,
consumerName: consumerName,
batchSize: 10,
blockMS: 5000,
}
}
// Setup 初始化消费者组(幂等操作)
func (c *Consumer) Setup(ctx context.Context) error {
// 尝试创建消费者组,从最新消息开始($)
// MKSTREAM 确保 Stream 存在时自动创建
err := c.client.XGroupCreateMkStream(ctx, c.stream, c.group, "$").Err()
if err != nil && err.Error() != "BUSYGROUP Consumer Group name already exists" {
return fmt.Errorf("create consumer group failed: %w", err)
}
log.Printf("[Consumer-%s] Group ready: %s", c.consumerName, c.group)
return nil
}
// ProcessEvents 消息处理主循环
func (c *Consumer) ProcessEvents(ctx context.Context) {
for {
select {
case <-ctx.Done():
log.Printf("[Consumer-%s] Shutting down...", c.consumerName)
return
default:
}
// XREADGROUP 读取消息
// ">" 表示只读取从未分配给任何消费者的新消息
streams, err := c.client.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: c.group,
Consumer: c.consumerName,
Streams: []string{c.stream, ">"},
Count: c.batchSize,
Block: time.Duration(c.blockMS) * time.Millisecond,
}).Result()
if err == redis.Nil {
// 阻塞超时,继续下一轮
continue
}
if err != nil {
log.Printf("[Consumer-%s] Read error: %v", c.consumerName, err)
time.Sleep(time.Second)
continue
}
for _, stream := range streams {
for _, message := range stream.Messages {
c.handleMessage(ctx, message)
}
}
}
}
// handleMessage 处理单条消息
func (c *Consumer) handleMessage(ctx context.Context, msg redis.XMessage) {
dataStr, ok := msg.Values["data"].(string)
if !ok {
log.Printf("[Consumer-%s] Invalid message format: %s", c.consumerName, msg.ID)
c.ack(ctx, msg.ID) // 格式错误也 ACK,避免阻塞
return
}
var event OrderEvent
if err := json.Unmarshal([]byte(dataStr), &event); err != nil {
log.Printf("[Consumer-%s] Unmarshal failed: %v", c.consumerName, err)
c.ack(ctx, msg.ID)
return
}
// 模拟业务处理(发送通知、更新统计等)
log.Printf("[Consumer-%s] Processing order=%s status=%s amount=%.2f",
c.consumerName, event.OrderID, event.Status, event.Amount)
// 模拟处理耗时
time.Sleep(50 * time.Millisecond)
// 处理成功,发送 ACK
if err := c.ack(ctx, msg.ID); err != nil {
log.Printf("[Consumer-%s] ACK failed for %s: %v", c.consumerName, msg.ID, err)
}
}
// ack 确认消息
func (c *Consumer) ack(ctx context.Context, ids ...string) error {
return c.client.XAck(ctx, c.stream, c.group, ids...).Err()
}
// ClaimPending 转移长时间未处理的待处理消息(故障恢复)
func (c *Consumer) ClaimPending(ctx context.Context, minIdleTime time.Duration) error {
// XAUTOCLAIM 自动转移 idle 超时的消息(Redis 6.2+)
// go-redis v9 支持 XAutoClaim
result, err := c.client.XAutoClaim(ctx, &redis.XAutoClaimArgs{
Stream: c.stream,
Group: c.group,
Consumer: c.consumerName,
MinIdle: minIdleTime,
Start: "0-0", // 从最小 ID 开始扫描
Count: 100,
}).Result()
if err != nil {
return err
}
for _, msg := range result.Messages {
log.Printf("[Consumer-%s] Claimed message %s from dead worker",
c.consumerName, msg.ID)
c.handleMessage(ctx, msg)
}
return nil
}
// Close 关闭连接
func (c *Consumer) Close() error {
return c.client.Close()
}
func main() {
// 通过环境变量区分不同 worker 实例
consumerName := os.Getenv("CONSUMER_NAME")
if consumerName == "" {
consumerName = "worker-1"
}
consumer := NewConsumer("localhost:6379", "order-processors", consumerName)
ctx, cancel := context.WithCancel(context.Background())
if err := consumer.Setup(ctx); err != nil {
log.Fatal(err)
}
// 启动消费者主循环
go consumer.ProcessEvents(ctx)
// 定期执行 ClaimPending(每 30 秒检查一次超时消息)
go func() {
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := consumer.ClaimPending(ctx, 60*time.Second); err != nil {
log.Printf("[Consumer-%s] Claim failed: %v", consumerName, err)
}
}
}
}()
// 优雅退出
sig := make(chan os.Signal, 1)
signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
<-sig
cancel()
consumer.Close()
log.Println("Consumer stopped gracefully")
}
5.5 部署与运行
# docker-compose.yml
version: "3.8"
services:
redis:
image: redis:7-alpine
ports:
- "6379:6379"
command: redis-server --appendonly yes
# 启动 Redis
docker-compose up -d
# 启动生产者
go run producer.go
# 启动多个消费者(不同终端)
CONSUMER_NAME=worker-1 go run consumer.go
CONSUMER_NAME=worker-2 go run consumer.go
CONSUMER_NAME=worker-3 go run consumer.go
5.6 PEL 监控脚本
#!/bin/bash
# monitor-pel.sh - 监控待处理消息
STREAM="orders:events"
GROUP="order-processors"
while true; do
PENDING=$(redis-cli XPENDING $STREAM $GROUP)
echo "[$(date '+%Y-%m-%d %H:%M:%S')] PEL Count: $PENDING"
if [ "$PENDING" -gt "0" ]; then
# 查看详情
redis-cli XPENDING $STREAM $GROUP - + 10
fi
sleep 5
done
六、最佳实践总结
Pub/Sub 只用于实时广播:在线通知、配置热更等允许消息丢失的场景。对于可靠性要求高的业务,一律使用 Streams。
Streams 必须配置 MAXLEN:防止内存无限增长,使用
~近似裁剪平衡性能与内存。ACK 是强要求:处理成功后立即
XACK,失败时记录日志后也建议 ACK 避免阻塞,通过监控发现失败再人工或自动补偿。Consumer Group 取名规范:按业务领域命名(如
order-processors、payment-notifiers),避免一个 Stream 挂载过多 Group。XAUTOCLAIM 做故障恢复:每个 Worker 定期执行,转移宕机同事遗留的待处理消息。MinIdle 建议设置为平均处理时间的 3-5 倍。
监控三个核心指标:PEL 堆积量、消费速率(msg/s)、ACK 延迟分布。任一指标异常都应触发告警。
连接池优化:消费者使用独立 Redis 连接(PoolSize 建议 5-10),生产者和消费者不共享连接池。
Redis Streams 在中小型消息队列场景下提供了"足够好"的解决方案:低延迟、低运维成本、与现有缓存基础设施复用。理解它的能力边界,才能做出正确的技术选型。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。