Redis 消息队列:Pub/Sub 广播与 Streams 持久化消息实战

对比 Redis Pub/Sub 和 Streams 两种消息队列方案,覆盖广播发布订阅、Consumer Group、ACK 确认、Kafka/RabbitMQ 选型对比

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.paidorders.shipped 等),PUBLISH 向频道广播消息。注意:所有命令在订阅状态下会阻塞 Redis 连接,实际应用中通常使用独立连接处理订阅,另一个连接处理其他 Redis 操作。

1.2 消息传递机制

Pub/Sub 的关键特征是推(Push)模型:

  1. 发布者调用 PUBLISH channel message
  2. Redis Server 遍历该频道的所有订阅者列表
  3. 将消息立即推送到每个订阅者的客户端缓冲区
  4. 如果订阅者网络延迟或处理缓慢,消息堆积在服务器输出缓冲区
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 范围查询
阻塞读取BRPOPXREAD BLOCK
内存效率中(附加元数据开销)

结论:所有新开发的消息队列场景都应使用 Streams,List 队列已进入维护模式。


四、Kafka / RabbitMQ / Redis Streams 选型对比

维度Redis StreamsApache KafkaRabbitMQ
数据持久化内存为主,RDB/AOF 落盘磁盘持久化,页缓存优化内存/磁盘可选
消息保留MAXLEN/~ 手动或自动裁剪时间/大小策略,长期保留TTL/队列长度限制
吞吐量中高(~100K msg/s 单机)极高(百万级/秒集群)中高(~50K msg/s)
延迟极低(亚毫秒)低(毫秒级)低(毫秒级)
消费者模型Push + PullPullPush + 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

六、最佳实践总结

  1. Pub/Sub 只用于实时广播:在线通知、配置热更等允许消息丢失的场景。对于可靠性要求高的业务,一律使用 Streams。

  2. Streams 必须配置 MAXLEN:防止内存无限增长,使用 ~ 近似裁剪平衡性能与内存。

  3. ACK 是强要求:处理成功后立即 XACK,失败时记录日志后也建议 ACK 避免阻塞,通过监控发现失败再人工或自动补偿。

  4. Consumer Group 取名规范:按业务领域命名(如 order-processorspayment-notifiers),避免一个 Stream 挂载过多 Group。

  5. XAUTOCLAIM 做故障恢复:每个 Worker 定期执行,转移宕机同事遗留的待处理消息。MinIdle 建议设置为平均处理时间的 3-5 倍。

  6. 监控三个核心指标:PEL 堆积量、消费速率(msg/s)、ACK 延迟分布。任一指标异常都应触发告警。

  7. 连接池优化:消费者使用独立 Redis 连接(PoolSize 建议 5-10),生产者和消费者不共享连接池。

Redis Streams 在中小型消息队列场景下提供了"足够好"的解决方案:低延迟、低运维成本、与现有缓存基础设施复用。理解它的能力边界,才能做出正确的技术选型。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. 缓存架构演进之路:从单机 Redis 到亿级分布式多级缓存体系
  2. Redis 7.x 重大新特性与架构升级深度解析
  3. Redis 消息队列深度对比:Pub/Sub、Streams 与 Kafka/RabbitMQ 选型指南