目录

Redis-09 发布订阅与消息队列

1. 发布订阅(Pub/Sub)

1.1 基本命令

# 订阅端
SUBSCRIBE channel [channel ...]           # 订阅一个或多个频道
PSUBSCRIBE pattern [pattern ...]          # 按模式订阅(glob 风格)
UNSUBSCRIBE [channel ...]                 # 取消订阅(不带参数取消全部)
PUNSUBSCRIBE [pattern ...]

# 发布端
PUBLISH channel message                   # 返回收到消息的订阅者数量

# 查询
PUBSUB CHANNELS [pattern]                 # 列出当前有订阅者的频道
PUBSUB NUMSUB [channel ...]               # 每个频道的订阅者数量
PUBSUB NUMPAT                             # 模式订阅的数量
PUBSUB SHARDCHANNELS [pattern]            # 7.0+ 分片频道
PUBSUB SHARDNUMSUB [channel ...]

# 7.0+ 集群分片 pub/sub
SSUBSCRIBE shardchannel
SPUBLISH shardchannel message
SUNSUBSCRIBE [shardchannel ...]

1.2 使用示例

开两个 redis-cli 窗口:

# 窗口 1:订阅
127.0.0.1:6379> SUBSCRIBE news:tech news:sport
Reading messages... (press Ctrl-C to quit)
1) "subscribe"
2) "news:tech"
3) (integer) 1        # 当前订阅数
1) "subscribe"
2) "news:sport"
3) (integer) 2

# 窗口 2:发布
127.0.0.1:6379> PUBLISH news:tech "Redis 8.0 released"
(integer) 1           # 有 1 个订阅者收到

# 窗口 1 立即收到
1) "message"
2) "news:tech"
3) "Redis 8.0 released"

模式订阅:

# 窗口 1
127.0.0.1:6379> PSUBSCRIBE news:*
1) "psubscribe"
2) "news:*"
3) (integer) 1

# 窗口 2
127.0.0.1:6379> PUBLISH news:tech "hello"
(integer) 1

# 窗口 1 收到(注意消息类型是 pmessage,多了一个 pattern 字段)
1) "pmessage"
2) "news:*"           # 匹配的模式
3) "news:tech"        # 实际的频道
4) "hello"

注意:如果一个客户端同时用 SUBSCRIBE news:tech 和 PSUBSCRIBE news:* 订阅,它会收到同一条消息两次(一次 message,一次 pmessage)。

1.3 实现原理

struct redisServer {
    dict *pubsub_channels;    // 频道 → 订阅它的客户端链表
    dict *pubsub_patterns;    // 模式 → 订阅它的客户端链表(7.0 前是 list)
    dict *pubsub_shardchannels; // 7.0+ 分片频道
};

SUBSCRIBE channel:

把 client 加入 server.pubsub_channels[channel] 这个链表
同时在 client->pubsub_channels 里也记一份(方便断开时清理)

PUBLISH channel msg:

1. 在 pubsub_channels 里查 channel,遍历订阅者链表,逐个 addReply 推送消息
2. 遍历 pubsub_patterns 里的所有模式,用 stringmatchlen() 逐个做 glob 匹配,
   匹配上的订阅者也推送
   ↑ 注意这是【遍历所有模式】,模式很多时 PUBLISH 会变慢(O(N*M))
3. 返回收到消息的订阅者总数

关键认识:PUBLISH 是一次"即时的、无状态的转发"。Redis 收到消息后立即推给当前在线的订阅者,然后就把消息丢掉了——不存储、不记录、不确认。

1.4 订阅模式下的限制

RESP2 协议下,进入订阅状态的连接不能再执行普通命令:

127.0.0.1:6379> SUBSCRIBE ch
Reading messages...
# 此时只能执行:SUBSCRIBE / UNSUBSCRIBE / PSUBSCRIBE / PUNSUBSCRIBE / PING / QUIT / RESET
127.0.0.1:6379> GET k
(error) ERR Can't execute 'get': only (P|S)SUBSCRIBE / (P|S)UNSUBSCRIBE / PING / QUIT / RESET are allowed in this context

为什么有这个限制:RESP2 里,服务端主动推送的消息和"命令的回复"用的是同一种编码(都是数组),客户端无法区分收到的数组是"我上一个命令的返回值"还是"服务端推来的消息"。

RESP3 解除了这个限制:它用 > 前缀明确标识 Push 类型消息,客户端能区分开,所以订阅连接可以同时执行普通命令:

127.0.0.1:6379> HELLO 3
127.0.0.1:6379> SUBSCRIBE ch
127.0.0.1:6379> GET k        # RESP3 下这是允许的

实践影响:RESP2 下客户端必须为订阅单独开一个连接,不能复用连接池里的连接。

1.5 Pub/Sub 的致命局限

这是面试必答的内容。

(1)消息不持久化,发出即忘(fire and forget)

PUBLISH 时不在线的订阅者,永久丢失这条消息。没有任何补偿机制。

时刻 T1:订阅者在线
时刻 T2:订阅者网络断开重连(可能只有 1 秒)
时刻 T3:PUBLISH 消息 M      ← 这条消息该订阅者永久收不到
时刻 T4:订阅者重连成功

(2)没有消息确认(ACK)机制

Redis 把消息 write 到订阅者的 socket 就认为完成了。订阅者收到后处理失败、或者收到后立刻崩溃,Redis 完全不知道,也不会重发。

(3)没有消费组,无法负载均衡

Pub/Sub 是广播语义:一条消息会推给所有订阅该频道的客户端。

如果你想"3 个 worker 分摊处理任务",用 Pub/Sub 会变成"3 个 worker 各自都处理了同一个任务"(重复消费)。想做负载均衡只能业务侧自己想办法(比如按消息内容哈希后各 worker 只处理自己那份,非常笨拙)。

(4)消息堆积会导致订阅者被断开

如果发布速度大于订阅者的消费速度,消息会堆积在该订阅者在服务端的输出缓冲区里。一旦超过限制:

client-output-buffer-limit pubsub 32mb 8mb 60
# hard limit 32MB:立即断开
# soft limit 8MB 持续 60 秒:断开

订阅者被断开 → 重连 → 期间的消息全部丢失。这是 Pub/Sub 在高流量场景下最常见的故障模式。

(5)PUBLISH 的消息也会占用主线程

PUBLISH 要遍历所有订阅者逐个推送,订阅者很多时(比如 1 万个)单次 PUBLISH 就是 1 万次 addReply。模式订阅更糟——要遍历所有模式做 glob 匹配。

(6)集群模式下的广播放大(7.0 前)

Redis Cluster 中,为了保证"订阅者连在任意节点都能收到消息",7.0 之前的 PUBLISH 会广播到集群里所有节点。这在大集群里造成严重的内部网络开销(一条消息在 10 节点集群里产生 10 倍流量)。

7.0 引入了分片 pub/sub(SPUBLISH/SSUBSCRIBE):分片频道名按 CRC16 算 slot,消息只在负责该 slot 的节点及其从节点之间传播,不再全集群广播。代价是订阅者必须连到正确的节点(客户端要处理 MOVED 重定向)。

1.6 Pub/Sub 的正确使用场景

明确了局限,就知道它只适合"允许丢消息"的场景:

场景 说明
配置热更新通知 通知所有应用实例"配置变了,去重新拉一次"。丢了也不致命(下次心跳会拉)
缓存失效广播 通知所有节点清本地缓存。可以配合定期全量刷新兜底
实时监控/日志推送 丢几条无所谓
在线状态、聊天室 IM 的"正在输入"提示这类瞬时状态
分布式锁的释放通知 Redisson 用 pub/sub 通知等锁的客户端"锁释放了,快来抢",即使丢了也有超时重试兜底
哨兵内部通信 Redis Sentinel 自己就用 pub/sub 在哨兵之间交换信息(__sentinel__:hello)

绝对不能用 Pub/Sub 的场景:订单处理、支付通知、任务分发、任何"消息不能丢"的业务。

1.7 一个典型用法:缓存失效广播

// 发布端:数据更新后通知所有节点
func UpdateUser(ctx context.Context, rdb *redis.Client, id int64, u *User) error {
    if err := db.Update(u); err != nil {
        return err
    }
    rdb.Del(ctx, fmt.Sprintf("user:%d", id))
    // 通知所有应用实例清理本地缓存
    rdb.Publish(ctx, "cache:invalidate", fmt.Sprintf("user:%d", id))
    return nil
}

// 订阅端:每个应用实例启动时开一个 goroutine
func WatchInvalidation(ctx context.Context, rdb *redis.Client, localCache *sync.Map) {
    // 注意:用独立连接,不占用连接池
    pubsub := rdb.Subscribe(ctx, "cache:invalidate")
    defer pubsub.Close()

    ch := pubsub.Channel()
    for {
        select {
        case msg, ok := <-ch:
            if !ok {
                return
            }
            localCache.Delete(msg.Payload)
        case <-ctx.Done():
            return
        }
    }
}

注意:因为 pub/sub 会丢消息,本地缓存必须同时设置较短的 TTL 作为兜底,不能完全依赖失效通知。

顺便一提:Redis 6.0 的**客户端缓存(Client-side Caching / CLIENT TRACKING)**就是这个模式的官方实现,而且更可靠——服务端会记录每个客户端缓存了哪些 key,精准推送 invalidate 消息。


2. 键空间通知(Keyspace Notification)

2.1 是什么

Redis 可以在数据发生变化时自动 PUBLISH 一条消息,让你能"监听某个 key 的变化"。

# 默认关闭(空字符串)
notify-keyspace-events ""

# 开启示例
notify-keyspace-events "KEA"      # 全部事件(最全,但开销大)
notify-keyspace-events "Ex"       # 只要 key 过期事件(最常用)
notify-keyspace-events "KEg$lshzxet"

2.2 配置字符的含义

字符 含义
K Keyspace 事件,频道格式 __keyspace@<db>__:<key>,消息内容是事件名
E Keyevent 事件,频道格式 __keyevent@<db>__:<event>,消息内容是key 名
g 通用命令(DEL、EXPIRE、RENAME 等)
$ string 命令
l list 命令
s set 命令
h hash 命令
z zset 命令
x 过期事件(key 过期时触发)
e 淘汰事件(key 被 maxmemory 策略淘汰时触发)
t stream 命令
d module key type 事件
m key miss 事件(6.0+,访问不存在的 key)
n 新 key 事件(7.0+)
A g$lshzxetd 的别名(不含 m 和 n)

必须包含 K 或 E 中至少一个,否则不会发出任何通知。

2.3 K 和 E 的区别

# 开启
127.0.0.1:6379> CONFIG SET notify-keyspace-events "KEA"

# 订阅端
127.0.0.1:6379> PSUBSCRIBE '__key*@0__:*'

# 另一个窗口执行
127.0.0.1:6379> SET foo bar

订阅端收到两条消息:

# Keyspace 通知(K):频道里带 key,内容是事件名
1) "pmessage"
2) "__key*@0__:*"
3) "__keyspace@0__:foo"       ← 频道包含 key 名
4) "set"                      ← 内容是事件名

# Keyevent 通知(E):频道里带事件名,内容是 key
1) "pmessage"
2) "__key*@0__:*"
3) "__keyevent@0__:set"       ← 频道包含事件名
4) "foo"                      ← 内容是 key 名

怎么选:

  • 关心某个特定 key 的所有变化 → 用 K:SUBSCRIBE __keyspace@0__:mykey;
  • 关心某类事件发生在哪些 key 上 → 用 E:SUBSCRIBE __keyevent@0__:expired(这是最常见的用法)。

2.4 最常见用途:监听 key 过期

127.0.0.1:6379> CONFIG SET notify-keyspace-events Ex
127.0.0.1:6379> SUBSCRIBE __keyevent@0__:expired

# 另一个窗口
127.0.0.1:6379> SET session:1001 data EX 5
# 5 秒后订阅端收到:
1) "message"
2) "__keyevent@0__:expired"
3) "session:1001"

看起来可以用来做延迟任务(设一个 TTL,到期后收到通知就执行任务)。但这个方案有严重缺陷,生产环境不能用:

缺陷一:过期事件的触发时机不准

第 6 篇讲过,Redis 的过期删除是惰性删除 + 定期删除。expired 事件是在key 真正被删除时才发出的,不是"TTL 到达的瞬间"。

所以如果一个 key 过期后:

  • 没人访问它(不触发惰性删除);
  • 定期删除的随机采样一直没抽到它;

那么 expired 事件可能延迟几秒、几分钟甚至更久才发出。做定时任务显然不可接受。

缺陷二:消息会丢(这是致命的)

expired 事件是通过 pub/sub 发送的,所以完全继承了 pub/sub 的所有缺陷:

  • 订阅者不在线(重启、部署、网络抖动)→ 消息永久丢失 → 任务永远不会执行;
  • 没有 ACK,处理失败不会重试;
  • 输出缓冲区满会断开连接。

缺陷三:广播导致重复消费

如果为了高可用部署了多个消费者实例,它们都会收到同一个过期事件,导致任务被重复执行。想去重又要额外加分布式锁。

缺陷四:主从环境下从节点收不到自己产生的事件

从节点不主动删除过期 key(等主节点发 DEL),所以:

  • 事件是在主节点上产生并发布的;
  • 连在从节点上的订阅者收不到 expired 事件(6.0 之前)。

Redis 6.0 做了改进:从节点在收到主节点的 DEL 后也会发出 expired 事件。但仍需注意版本差异。

结论:键空间通知只适合做"辅助性的、可丢失的"监听(比如统计、审计、缓存联动),绝不能作为业务逻辑的唯一触发源。延迟任务应该用 zset 轮询或专业的延迟队列(第 3 节)。

2.5 性能开销

开启 notify-keyspace-events "KEA" 会让每个写命令都额外产生 1~2 次 PUBLISH。即使没有任何订阅者,Redis 也要走一遍 pubsubPublishMessage 的逻辑(查字典、遍历模式)。

实测在高 QPS 下,开启全量通知可能带来 10%~30% 的吞吐下降。

建议:只开真正需要的事件类型(如 Ex 只要过期事件),绝不要在生产开 KEA。


3. 用 List 做消息队列

在 Stream 出现之前(5.0 前),list 是 Redis 做队列的主流方案。理解它的演进过程很有价值。

3.1 方案一:LPUSH + RPOP(轮询)

# 生产者
LPUSH queue:task '{"id":1,"type":"email"}'

# 消费者:轮询
while true:
    task = RPOP queue:task
    if task: process(task)
    else: sleep(1)      # 队列空,睡一会再试

问题:

  • 睡眠时间难权衡:睡得久 → 消息处理延迟高;睡得短 → 大量无效的 RPOP 请求浪费 CPU 和网络(假设 100 个消费者每 100ms 轮询一次,就是 1000 QPS 的空请求)。

3.2 方案二:BRPOP(阻塞)

# 消费者:阻塞等待,有消息立即返回
BRPOP queue:task 0      # 0 表示永久阻塞

解决了轮询问题(第 5 篇讲了阻塞不占主线程、按 FIFO 唤醒)。

但仍有严重问题:消息可能丢失

1. 消费者 BRPOP 拿到消息 M(此时 M 已从 list 中删除)
2. 消费者正在处理 M 时【崩溃】
3. 消息 M 永久丢失——Redis 里没有它,也没人知道它没被处理完

3.3 方案三:BRPOPLPUSH / BLMOVE(可靠队列)

# 消费者:从任务队列取出的同时,原子地放进"处理中"队列
BRPOPLPUSH queue:task queue:processing 0
# 6.2+ 推荐用更灵活的 BLMOVE
BLMOVE queue:task queue:processing RIGHT LEFT 0

# 处理成功后从 processing 中删除
LREM queue:processing 1 '<消息内容>'

这是 Redis 官方文档里叫做 “Reliable queue”(可靠队列) 的模式。

配套需要一个"看门狗"进程:

# 定期检查 processing 队列里滞留太久的消息,重新投递
# 问题:list 里没有时间信息,无法知道消息在 processing 里待了多久!

这就暴露了这个方案的根本缺陷:

  1. 无法知道消息何时进入 processing(list 元素没有时间戳),只能:
    • 在消息体里自己塞时间戳(消费者取出后要修改消息内容再放进 processing,破坏了原子性);
    • 或者用一个额外的 zset 记录 消息 → 进入时间(复杂度上升,且两个操作的原子性又需要 Lua);
  2. LREM 是 O(N):processing 队列长时要遍历查找,性能差;
  3. 无法区分是哪个消费者在处理(排查困难);
  4. 仍然没有消费组:多个业务系统想各自完整消费一遍做不到;
  5. 无法回溯:消息处理完就删了,想重放不可能。

3.4 list 队列的其他固有缺陷

(1)无法广播 / 无消费组

一条消息只能被一个消费者拿到。如果订单创建后既要发短信、又要加积分、又要写数仓,三个系统各自需要一份消息,list 做不到(只能生产者写三个 list,耦合严重)。

(2)无法回溯与重放

消息弹出即删除。生产事故后想"重放最近一小时的消息"不可能。

(3)无法查看队列状态

不知道每个消费者处理了多少、有没有卡住、消费延迟多大。只能 LLEN 看积压总量。

(4)阻塞消费者独占连接

BRPOP 期间该连接不能复用,N 个消费者要 N 个连接。

(5)优先级和延迟做不到

要延迟就得用 zset,要优先级得开多个 list(BRPOP high mid low 0 勉强能做,但只有粗粒度的三档)。

3.5 结论

如果你的 Redis 版本 >= 5.0,做消息队列请直接用 Stream,不要再用 list。

list 只在这两种情况下还合适:

  1. 极简的、允许丢失的任务队列(比如异步发送非关键通知);
  2. 只需要"最新 N 条"的场景(LPUSH + LTRIM,这不算队列)。

4. 用 ZSet 做延迟队列

延迟队列是 list 完全做不到、而 Stream 也不直接支持的场景,zset 是最常见的方案。

4.1 基本思路

用 score 存"应该执行的时间戳",消费者轮询取出所有 score <= 当前时间的成员。

# 生产:30 秒后执行
ZADD delay:queue <当前时间戳+30000> '{"id":1,"task":"cancel_order","orderId":1001}'

# 消费:取已到期的任务
ZRANGEBYSCORE delay:queue 0 <当前时间戳> LIMIT 0 10

4.2 必须用 Lua 保证原子性

如果分成"查询"和"删除"两步,多个消费者会重复消费同一个任务:

消费者 A: ZRANGEBYSCORE → 拿到 task1
消费者 B: ZRANGEBYSCORE → 也拿到 task1   ← 重复!
消费者 A: ZREM task1
消费者 B: ZREM task1(返回 0,但已经处理了两次)

正确做法:

-- KEYS[1]=延迟队列key  ARGV[1]=当前时间戳  ARGV[2]=一次取多少
local tasks = redis.call('ZRANGEBYSCORE', KEYS[1], 0, ARGV[1], 'LIMIT', 0, ARGV[2])
if #tasks > 0 then
    redis.call('ZREM', KEYS[1], unpack(tasks))
end
return tasks

注意 unpack 有栈限制(约 8000 个参数),所以 ARGV[2] 要控制在几百以内。

4.3 完整实现(Go)

type DelayQueue struct {
    rdb  *redis.Client
    key  string
}

var fetchScript = redis.NewScript(`
    local tasks = redis.call('ZRANGEBYSCORE', KEYS[1], 0, ARGV[1], 'LIMIT', 0, ARGV[2])
    if #tasks > 0 then
        redis.call('ZREM', KEYS[1], unpack(tasks))
    end
    return tasks
`)

// 投递延迟任务
func (q *DelayQueue) Push(ctx context.Context, payload string, delay time.Duration) error {
    executeAt := time.Now().Add(delay).UnixMilli()
    return q.rdb.ZAdd(ctx, q.key, redis.Z{
        Score:  float64(executeAt),
        Member: payload,
    }).Err()
}

// 消费循环
func (q *DelayQueue) Consume(ctx context.Context, handler func(string) error) {
    ticker := time.NewTicker(200 * time.Millisecond)   // 轮询间隔
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            now := time.Now().UnixMilli()
            res, err := fetchScript.Run(ctx, q.rdb,
                []string{q.key}, now, 100).StringSlice()
            if err != nil {
                log.Printf("fetch error: %v", err)
                continue
            }
            for _, payload := range res {
                if err := handler(payload); err != nil {
                    // 处理失败:重新投递(可以加退避)
                    q.Push(ctx, payload, 10*time.Second)
                }
            }
        }
    }
}

4.4 zset 延迟队列的问题与改进

问题一:轮询有延迟,且空轮询浪费资源

  • 轮询间隔 200ms 意味着任务最多延迟 200ms 执行;
  • 大部分轮询是空的(没有到期任务),浪费请求。

改进:动态调整轮询间隔——用 ZRANGE key 0 0 WITHSCORES 看最近的任务什么时候到期,如果还有 5 秒,就睡 5 秒(但要能被新任务打断,可以用 channel + 定时器)。

问题二:任务丢失(取出后消费者崩溃)

Lua 脚本"取出即删除",取出后崩溃任务就丢了。

改进:两阶段确认——取出时不删除,而是把 score 改成"当前时间 + 处理超时"(相当于租约):

-- 取出并"续期"(把 score 改到未来,其他消费者暂时看不到它)
local tasks = redis.call('ZRANGEBYSCORE', KEYS[1], 0, ARGV[1], 'LIMIT', 0, ARGV[2])
for i, task in ipairs(tasks) do
    -- 租约 60 秒:如果 60 秒内没 ACK(ZREM),任务会重新变成"到期"被别人取走
    redis.call('ZADD', KEYS[1], tonumber(ARGV[1]) + 60000, task)
end
return tasks

处理成功后再 ZREM 确认删除。这就实现了简易的"至少一次"投递。

问题三:无法重复投递相同内容的任务

zset 的 member 唯一。如果两次投递完全相同的 payload,第二次只会更新 score而不是新增一条。

改进:给 payload 加唯一 ID({"uuid":"xxx", "data":...}),或者用 taskId 作为 member、任务详情存在另一个 hash 里。

问题四:大 key

所有延迟任务都在一个 zset 里,任务多时形成大 key(ZADD/ZRANGEBYSCORE 变慢、迁移困难)。

改进:按时间分片——比如按小时分 key(delay:queue:2026072913),消费者同时轮询当前小时和上一小时的 key。或者按业务类型分 key。

4.5 延迟队列的其他实现方式

方案 优点 缺点
zset 轮询(上面的) 简单、精度高(毫秒级) 需要自己实现可靠性、轮询开销
多级时间轮 + list 内存效率高 实现复杂
键空间过期通知 不需要轮询 不可靠(会丢消息)+ 时机不准,不要用
RocketMQ / RabbitMQ 延迟消息 生产级可靠 需要额外中间件
数据库轮询 + 索引 可靠、可查询 性能差、给 DB 压力

结论:Redis 做延迟队列首选 zset + Lua + 租约机制 + 时间分片。如果业务对可靠性要求极高(比如订单超时关闭),建议双保险:Redis 延迟队列做主路径(快),数据库定时任务扫描做兜底(慢但可靠)。


5. Stream 消息队列完整实践

第 3 篇介绍了 Stream 的命令,这里讲工程实践。

5.1 核心概念回顾

Stream(append-only 日志)
├── 消息:ID(毫秒时间戳-序号)+ field-value 对
├── 消费组 A
│   ├── last-delivered-id:这个组投递到哪了
│   ├── PEL(Pending Entries List):已投递未确认的消息
│   ├── 消费者 c1(有自己的 PEL 子集)
│   └── 消费者 c2
└── 消费组 B(独立的消费进度,互不影响)

关键点:

  • 消息持久化在 stream 里,不会因为被消费而删除(除非 XDEL/XTRIM);
  • 每个消费组独立维护进度(last-delivered-id),实现"多个业务系统各自完整消费一遍";
  • 组内多个消费者分摊消息(一条消息只投给组内一个消费者),实现负载均衡;
  • PEL 记录未确认消息,支持崩溃恢复和重投。

5.2 生产者最佳实践

func Produce(ctx context.Context, rdb *redis.Client, order *Order) error {
    id, err := rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "stream:orders",
        // 关键:一定要设置裁剪策略,否则内存无限增长
        MaxLen: 100000,          // 保留最近 10 万条
        Approx: true,            // 用 ~ 近似裁剪,性能好得多
        ID:     "*",             // 让 Redis 生成 ID
        Values: map[string]interface{}{
            "orderId": order.ID,
            "userId":  order.UserID,
            "amount":  order.Amount,
            "ts":      time.Now().UnixMilli(),
        },
    }).Result()
    if err != nil {
        return err
    }
    log.Printf("produced message %s", id)
    return nil
}

要点:

  1. 必须设置 MaxLen 或 MinID 裁剪,否则内存会无限增长;
  2. 用 Approx: true(对应 ~)而不是精确裁剪:近似裁剪只在能整个删掉一个 listpack 节点时才删(O(1) 摊还),精确裁剪要逐条删(可能 O(N) 阻塞);
  3. 裁剪阈值要留足余量:如果消费组有滞后(lag),MaxLen 太小会直接丢掉未消费的消息。要监控 XINFO GROUPS 的 lag;
  4. 字段名要短且固定:Stream 的 listpack 对同一节点内重复的字段名做了去重优化(master entry 机制),字段名固定能省大量内存。

5.3 消费者完整实现

一个生产级的消费者要处理三件事:正常消费、崩溃恢复、超时消息认领。

type StreamConsumer struct {
    rdb      *redis.Client
    stream   string
    group    string
    consumer string   // 每个实例一个唯一名字,如 "worker-hostname-pid"
    handler  func(msg redis.XMessage) error
}

func (c *StreamConsumer) Start(ctx context.Context) error {
    // ---------- 步骤 1:确保消费组存在 ----------
    // MKSTREAM:stream 不存在时自动创建
    // "0" 表示从头消费;"$" 表示只消费此后的新消息
    err := c.rdb.XGroupCreateMkStream(ctx, c.stream, c.group, "0").Err()
    if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") {
        return err   // BUSYGROUP 表示组已存在,可忽略
    }

    // ---------- 步骤 2:崩溃恢复:先处理自己之前未 ACK 的消息 ----------
    // 用 "0" 而不是 ">",表示"给我这个消费者已投递但未确认的消息"
    if err := c.recoverPending(ctx); err != nil {
        log.Printf("recover error: %v", err)
    }

    // ---------- 步骤 3:启动超时消息认领的巡检协程 ----------
    go c.claimLoop(ctx)

    // ---------- 步骤 4:正常消费循环 ----------
    return c.consumeLoop(ctx)
}

// 处理自己的 PEL(重启后调用)
func (c *StreamConsumer) recoverPending(ctx context.Context) error {
    start := "0"
    for {
        res, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    c.group,
            Consumer: c.consumer,
            Streams:  []string{c.stream, start},   // 注意这里是 "0" 不是 ">"
            Count:    100,
        }).Result()
        if err != nil {
            if err == redis.Nil { return nil }
            return err
        }
        if len(res) == 0 || len(res[0].Messages) == 0 {
            return nil   // PEL 已清空
        }
        for _, msg := range res[0].Messages {
            c.process(ctx, msg)
            start = msg.ID   // 继续从这条之后读
        }
    }
}

// 正常消费
func (c *StreamConsumer) consumeLoop(ctx context.Context) error {
    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        default:
        }

        res, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    c.group,
            Consumer: c.consumer,
            Streams:  []string{c.stream, ">"},    // ">" = 从未投递过的新消息
            Count:    10,
            Block:    5 * time.Second,            // 阻塞等待,避免空轮询
        }).Result()

        if err != nil {
            if err == redis.Nil {
                continue     // Block 超时,没有新消息,正常
            }
            log.Printf("XReadGroup error: %v", err)
            time.Sleep(time.Second)
            continue
        }

        for _, stream := range res {
            for _, msg := range stream.Messages {
                c.process(ctx, msg)
            }
        }
    }
}

// 处理单条消息 + ACK
func (c *StreamConsumer) process(ctx context.Context, msg redis.XMessage) {
    if err := c.handler(msg); err != nil {
        log.Printf("handle message %s failed: %v", msg.ID, err)
        // 【不 ACK】,消息留在 PEL 里等重试
        return
    }
    // 处理成功才 ACK
    if err := c.rdb.XAck(ctx, c.stream, c.group, msg.ID).Err(); err != nil {
        log.Printf("ack %s failed: %v", msg.ID, err)
    }
}

// 认领其他消费者卡住的消息
func (c *StreamConsumer) claimLoop(ctx context.Context) {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            start := "0-0"
            for {
                // XAUTOCLAIM:认领空闲超过 60 秒的消息
                msgs, next, err := c.rdb.XAutoClaim(ctx, &redis.XAutoClaimArgs{
                    Stream:   c.stream,
                    Group:    c.group,
                    Consumer: c.consumer,
                    MinIdle:  60 * time.Second,
                    Start:    start,
                    Count:    10,
                }).Result()
                if err != nil {
                    log.Printf("autoclaim error: %v", err)
                    break
                }
                if len(msgs) == 0 {
                    break
                }
                for _, msg := range msgs {
                    c.processWithDeadLetter(ctx, msg)
                }
                if next == "0-0" {
                    break     // 遍历完一轮
                }
                start = next
            }
        }
    }
}

// 带死信处理的消费(用于重投的消息)
func (c *StreamConsumer) processWithDeadLetter(ctx context.Context, msg redis.XMessage) {
    // 查这条消息被投递了几次
    pending, err := c.rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream: c.stream,
        Group:  c.group,
        Start:  msg.ID,
        End:    msg.ID,
        Count:  1,
    }).Result()

    if err == nil && len(pending) > 0 && pending[0].RetryCount > 5 {
        // 【死信处理】:重试超过 5 次,说明是毒消息
        log.Printf("message %s exceeded max retries, moving to dead letter", msg.ID)
        c.rdb.XAdd(ctx, &redis.XAddArgs{
            Stream: c.stream + ":dead",
            Values: msg.Values,
        })
        c.rdb.XAck(ctx, c.stream, c.group, msg.ID)   // ACK 掉,不再重投
        return
    }

    c.process(ctx, msg)
}

5.4 关键设计点解析

(1)为什么启动时要先用 STREAMS key 0 读一遍

因为消费者重启后,它之前"已投递但未 ACK"的消息还在 PEL 里,用 > 是读不到的(> 只给"从未投递过"的新消息)。必须用 0(或具体 ID)把自己的 PEL 读出来重新处理。

如果不做这一步,那些消息会一直卡在 PEL 里,只能靠 XAUTOCLAIM 的巡检兜底(延迟至少 MinIdle 时间)。

(2)消费者名字必须稳定且唯一

  • 唯一:两个实例用同一个消费者名,PEL 会混在一起,XAUTOCLAIM 逻辑会错乱;
  • 稳定:如果每次重启都用随机名(如 UUID),旧名字的 PEL 就成了"孤儿"——那个消费者永远不会回来读自己的 PEL,只能等 XAUTOCLAIM。而且 XINFO CONSUMERS 里会堆积大量僵尸消费者。

推荐:用 hostname + 固定序号(如 K8s StatefulSet 的 pod 名 worker-0、worker-1),保证重启后名字不变。用 XGROUP DELCONSUMER 清理确定不再回来的消费者。

(3)XAUTOCLAIM 的 MinIdle 怎么设

必须 > 单条消息的最大处理时间。如果消息处理要 30 秒而 MinIdle 设 10 秒,那么消息处理中就被别人认领了,导致重复消费。

建议:MinIdle = 最大处理时间 × 2~3。

(4)死信处理是必须的

如果一条消息因为数据问题(比如 JSON 格式错误、关联数据被删)永远处理失败,不做死信处理它会被无限重投,浪费资源并可能阻塞后面的消息。

标准做法:RetryCount 超过阈值 → 转存到死信 stream → XACK 原消息。

5.5 监控指标

# stream 整体状况
127.0.0.1:6379> XINFO STREAM stream:orders
 1) "length"                        # 消息总数
 2) (integer) 100000
 3) "radix-tree-keys"               # radix tree 节点数
 4) (integer) 200
 5) "radix-tree-nodes"
 6) (integer) 450
 7) "last-generated-id"             # 最新消息 ID
 8) "1785657600123-0"
 9) "max-deleted-entry-id"          # 7.0+ 被删除的最大 ID
10) "0-0"
11) "entries-added"                 # 7.0+ 累计添加的消息数
12) (integer) 523400
13) "groups"                        # 消费组数量
14) (integer) 2
15) "first-entry" ...
16) "last-entry" ...

# 消费组状况(最重要)
127.0.0.1:6379> XINFO GROUPS stream:orders
1)  1) "name"
    2) "group:sms"
    3) "consumers"
    4) (integer) 3           # 消费者数量
    5) "pending"
    6) (integer) 12          # 【PEL 大小】未确认消息数
    7) "last-delivered-id"
    8) "1785657600100-0"
    9) "entries-read"        # 7.0+ 该组已读取的消息数
   10) (integer) 523388
   11) "lag"                 # 【最重要】还有多少消息没消费
   12) (integer) 12

# 消费者状况
127.0.0.1:6379> XINFO CONSUMERS stream:orders group:sms
1) 1) "name"
   2) "worker-0"
   3) "pending"
   4) (integer) 4
   5) "idle"                 # 距上次交互多久(毫秒)
   6) (integer) 1523
   7) "inactive"             # 7.2+ 距上次成功读取消息多久
   8) (integer) 1523

# PEL 明细
127.0.0.1:6379> XPENDING stream:orders group:sms - + 10
1) 1) "1785657600123-0"      # 消息 ID
   2) "worker-0"             # 归属消费者
   3) (integer) 300000       # 空闲毫秒数
   4) (integer) 3            # 【投递次数】超过阈值要做死信处理

必须告警的指标:

指标 阈值 含义
lag 持续增长 消费能力不足,需要加消费者或优化处理逻辑
pending 持续增长 大量消息未 ACK:处理失败、或者忘了 ACK、或者消费者崩溃
XPENDING 的投递次数 > 5 有毒消息在无限重试
消费者 idle 异常大 消费者卡死或已崩溃
length 接近 MaxLen 裁剪在生效,检查是否会丢未消费的消息

5.6 Stream 的运维注意点

(1)XDEL 不释放内存

XDEL 只是给 listpack 里的条目打删除标记,XLEN 会减少但内存不一定降。真正回收要靠 XTRIM(整个 listpack 节点被裁掉时才释放)。

(2)消费组的 PEL 会阻止裁剪的安全性

XTRIM MAXLEN 1000 会无条件删掉超出的消息,即使它们还在某个消费组的 PEL 里或者还没被消费。所以:

  • 裁剪阈值要远大于最大可能的 lag;
  • 监控 lag,发现消费严重滞后要先解决消费问题再裁剪。

(3)消费组会随 stream 一起被 DEL 删除

DEL stream:orders 会连消费组和 PEL 一起删掉,所有进度丢失。要小心。

(4)集群模式下一个 stream 在一个节点上

Stream 是单个 key,所以整个 stream 的数据在一个 slot/节点上,无法像 Kafka 那样分区并行。要提高吞吐只能自己做分片(stream:orders:0、stream:orders:1,生产者按 hash 选一个),但这样就丢失了全局顺序。

这是 Stream 相比 Kafka 最本质的扩展性限制。


6. 三种方案的完整对比

6.1 Redis 内部三种方案

能力 Pub/Sub List Stream
消息持久化 ❌ ✅ ✅
消费者离线不丢 ❌ ✅ ✅
消息确认(ACK) ❌ ❌ ✅
消息重投 ❌ 需自己实现(很难) ✅ XCLAIM/XAUTOCLAIM
消费组(多组独立消费) ❌(广播但不持久) ❌ ✅
组内负载均衡 ❌ 伪(抢) ✅
消息回溯/重放 ❌ ❌ ✅ XRANGE
阻塞消费 ✅ ✅ BRPOP ✅ BLOCK
消息 ID / 时间定位 ❌ ❌ ✅
多字段消息 ❌(要自己序列化) ❌ ✅
延迟消息 ❌ ❌ ❌(要用 zset)
消费进度可观测 ❌ 只有 LLEN ✅ XINFO(含 lag)
死信队列 ❌ 需自己实现 ✅ 靠 RetryCount + 死信 stream
实现复杂度 最简单 简单 中等

结论:

  • Pub/Sub:只用于"允许丢失的实时通知"(配置更新、缓存失效、在线状态);
  • List:极简的、允许丢失的任务队列,或"最新 N 条"场景;
  • Stream:做消息队列的唯一正确选择(5.0+)。

6.2 Redis Stream vs Kafka / RocketMQ

维度 Redis Stream Kafka
存储介质 内存(堆积就是内存成本) 磁盘顺序写(堆积几乎无成本)
存储上限 受内存限制(GB 级) TB ~ PB 级
消息保留 几小时~几天(内存扛不住更久) 几周~几个月甚至永久
吞吐量 10 万级/秒(单节点) 百万级/秒(多分区并行)
延迟 微秒~毫秒级(更低) 毫秒~十毫秒级
水平扩展 ❌ 单 stream 在单节点,要自己分片 ✅ 原生分区(partition)模型
消息顺序 全局有序(单 stream) 分区内有序
持久化保证 AOF everysec 可能丢 1 秒;异步复制会丢 ISR + acks=all 可保证不丢
Exactly-Once ❌ 只有 At-Least-Once ✅ 事务 + 幂等生产者
消费者组 rebalance 手动(XAUTOCLAIM) ✅ 自动 rebalance
生态 简单 Kafka Connect、Streams、Flink 集成、Schema Registry
运维成本 低(已有 Redis 就不用加组件) 高(ZK/KRaft + broker 集群)
适合场景 轻量任务队列、实时通知、缓冲削峰 数据管道、日志收集、流处理、事件溯源

6.3 选型决策

选 Redis Stream:

  • 消息量中等(每天千万级以内),堆积可控(几百 MB 到几 GB);
  • 消息生命周期短(几小时到几天);
  • 已经在用 Redis,不想为了一个队列引入 Kafka 的运维成本;
  • 需要极低延迟(微秒级);
  • 典型场景:异步任务(发短信/邮件)、订单状态流转、实时通知、削峰缓冲。

必须选 Kafka / RocketMQ:

  • 数据量 TB 级,需要长期保存和任意重放;
  • 需要百万级 TPS 和水平扩展(多分区并行消费);
  • 需要严格的不丢消息保证(Redis 的 AOF + 异步复制做不到);
  • 需要 Exactly-Once 语义、事务消息;
  • 需要流处理生态(Flink/Kafka Streams);
  • 典型场景:用户行为日志、数据同步管道、大数据 ETL、金融交易流水。

一句话:Stream 是"够用的轻量 MQ",不是 Kafka 的替代品。 用它的前提是接受"数据在内存里"这个根本约束。


7. 高频面试题

Q1:Redis 的 Pub/Sub 有什么局限?为什么不能用它做消息队列?

六个致命局限:

  1. 消息不持久化,发出即忘:PUBLISH 时不在线的订阅者永久丢失消息,没有任何补偿机制。订阅者哪怕只断开 1 秒,这期间的消息全丢;
  2. 没有 ACK 机制:Redis 把消息 write 到 socket 就认为完成,订阅者处理失败或立刻崩溃,Redis 完全不知道,也不会重发;
  3. 没有消费组,无法负载均衡:Pub/Sub 是广播语义,一条消息推给所有订阅者。想让 3 个 worker 分摊任务会变成 3 个 worker 各处理一遍;
  4. 堆积会导致订阅者被断开:消费慢时消息堆在服务端该订阅者的输出缓冲区,超过 client-output-buffer-limit pubsub 32mb 8mb 60 就被强制断开,期间消息全丢;
  5. PUBLISH 占用主线程:要遍历所有订阅者逐个推送;模式订阅更糟——要遍历所有模式做 glob 匹配(O(N*M));
  6. 集群下广播放大(7.0 前):为保证订阅者连任意节点都能收到,PUBLISH 会广播到所有节点,10 节点集群产生 10 倍流量。7.0 的分片 pub/sub(SPUBLISH/SSUBSCRIBE)解决了这点。

只适合"允许丢消息"的场景:配置热更新、缓存失效广播、实时监控推送、在线状态、分布式锁的释放通知(有超时兜底)。Redis Sentinel 内部通信也用它。

Q2:为什么订阅状态下不能执行其他命令?

RESP2 协议下,服务端主动推送的消息和"命令的回复"用同一种编码(都是数组),客户端无法区分收到的数组是"我上个命令的返回值"还是"服务端推来的消息"。所以 Redis 强制订阅连接只能执行 SUBSCRIBE/UNSUBSCRIBE/PSUBSCRIBE/PUNSUBSCRIBE/PING/QUIT/RESET。

RESP3(6.0+)解除了这个限制:它引入了 > 前缀的 Push 类型明确标识主动推送,客户端能区分,所以订阅连接可以同时执行普通命令(发 HELLO 3 协商)。

实践影响:RESP2 下客户端必须为订阅单独开连接,不能用连接池里的连接(否则那个连接就被占死了)。

Q3:能用 key 过期通知做延迟任务吗?

能实现但生产环境绝对不能用,四个致命问题:

  1. 触发时机严重不准:Redis 的过期删除是惰性 + 定期的,expired 事件在 key 真正被删除时才发出,不是 TTL 到达的瞬间。如果这个 key 没人访问、定期删除的随机采样也没抽到它,事件可能延迟几秒、几分钟甚至更久;
  2. 消息会丢(致命):expired 事件是通过 pub/sub 发送的,完全继承其缺陷——消费者不在线(重启/部署/网络抖动)→ 事件永久丢失 → 任务永远不会执行;
  3. 广播导致重复消费:部署多个消费者实例做高可用时,它们都会收到同一个事件,任务被重复执行,还得额外加分布式锁去重;
  4. 主从环境的差异:事件在主节点产生,6.0 之前连在从节点的订阅者收不到。

正确做法:延迟任务用 zset + Lua + 租约机制(score 存执行时间戳)。对可靠性要求极高的(如订单超时关闭)应该双保险:Redis 延迟队列走快路径 + 数据库定时任务扫描兜底。

另外:开启 notify-keyspace-events "KEA" 会让每个写命令额外产生 1~2 次 PUBLISH,高 QPS 下可能带来 10%~30% 的吞吐下降。生产只开真正需要的类型(如 Ex)。

Q4:用 list 做消息队列有什么问题?

方案演进与问题:

  1. LPUSH + RPOP 轮询:睡眠间隔难权衡——睡久了延迟高,睡短了大量空请求浪费资源;
  2. BRPOP 阻塞:解决了轮询问题,但消息会丢——BRPOP 拿到消息时它已从 list 删除,消费者处理中崩溃则消息永久丢失;
  3. BRPOPLPUSH/BLMOVE 可靠队列:取出的同时放进"处理中"队列,处理完 LREM 删除。但根本缺陷是:
    • list 元素没有时间信息,无法知道消息在 processing 里待了多久,看门狗进程没法判断超时;
    • LREM 是 O(N),processing 长时性能差;
    • 不知道是哪个消费者在处理,排查困难。

list 的固有缺陷(无法通过技巧绕过):

  • 无消费组:一条消息只能被一个消费者拿到,多个业务系统各自完整消费一遍做不到;
  • 无法回溯重放:弹出即删除,事故后想重放不可能;
  • 无消费进度可观测:只能 LLEN 看总积压,不知道每个消费者的状态;
  • 阻塞消费者独占连接;
  • 无延迟、无优先级(BRPOP q1 q2 q3 只能做粗粒度的三档优先级)。

结论:Redis >= 5.0 做消息队列请直接用 Stream。

Q5:Stream 的 PEL 是什么?消息不 ACK 会怎样?

PEL(Pending Entries List) 是每个消费组维护的"已投递但未确认"消息列表,记录每条消息的:ID、当前归属的消费者、最后投递时间、投递次数(RetryCount)。

消息不 XACK 的后果:

  1. 永远留在 PEL 里,XPENDING 数量持续增长;
  2. 不会自动重投——必须有人主动 XCLAIM/XAUTOCLAIM(按 min-idle-time 筛超时的),或者该消费者用 XREADGROUP ... STREAMS key 0 重新拉自己的 PEL;
  3. 消息本体也不能被安全裁剪(虽然 XTRIM 会无条件删,但那就丢消息了)。

所以生产上必须有两个机制:

  1. 巡检任务周期性执行 XAUTOCLAIM,认领崩溃消费者卡住的消息;
  2. 死信处理:RetryCount 超过阈值(如 5)的消息,说明是毒消息(数据格式错误、关联数据被删),要转存到死信 stream 并 XACK 掉,否则会被无限重投,浪费资源。

Q6:Stream 消费者重启后,之前未处理完的消息怎么办?

必须主动去读自己的 PEL:

# ">" 只给"从未投递过"的新消息,读不到 PEL 里的
XREADGROUP GROUP g1 consumer1 STREAMS orders >

# 用 "0"(或具体 ID)才能读到"已投递给我但未 ACK"的消息
XREADGROUP GROUP g1 consumer1 STREAMS orders 0

标准的消费者启动流程:

1. XGROUP CREATE ... MKSTREAM(忽略 BUSYGROUP 错误)
2. 【崩溃恢复】用 STREAMS key 0 循环读自己的 PEL,处理并 ACK,直到返回空
3. 启动 XAUTOCLAIM 巡检协程(认领其他消费者卡住的消息)
4. 进入正常消费循环:STREAMS key > + BLOCK

这也是"消费者名字必须稳定"的原因:如果每次重启用随机名(UUID),旧名字的 PEL 就成了孤儿——那个消费者永远不会回来读它,只能等 XAUTOCLAIM 兜底(延迟至少 MinIdle),而且 XINFO CONSUMERS 会堆积僵尸消费者。推荐用 K8s StatefulSet 的 pod 名(worker-0、worker-1)这种稳定标识。

Q7:XAUTOCLAIM 的 min-idle-time 怎么设?设错了会怎样?

必须大于"单条消息的最大处理时间"。

如果消息处理需要 30 秒,而 min-idle-time 设成 10 秒:消费者 A 正在处理消息(已经 15 秒),巡检发现它"空闲超过 10 秒"就把消息认领给了消费者 B → 同一条消息被两个消费者同时处理 → 重复消费(如果业务不幂等就出问题)。

建议:min-idle-time = 最大处理时间 × 2~3。

同时业务处理逻辑应该做幂等(用消息 ID 做去重键),因为 Stream 只提供 At-Least-Once 语义,重复投递在异常情况下总是可能发生的。

Q8:Stream 会无限增长吗?怎么控制?

会。Stream 是 append-only 的,不裁剪会一直涨到打爆内存。

控制手段:

  1. XADD 时带裁剪(推荐):XADD s MAXLEN ~ 100000 * f v 或 MINID ~ <时间戳>;
  2. 用 ~(近似裁剪)而不是 =(精确裁剪):近似裁剪只在能整个删掉一个 listpack 节点时才删(O(1) 摊还),精确裁剪要逐条删可能 O(N) 阻塞主线程;
  3. 定期 XTRIM:XTRIM s MINID ~ <7天前时间戳>。

两个危险点:

  1. XDEL 不释放内存:只是给 listpack 条目打删除标记,XLEN 减少但内存不降,真正回收靠 XTRIM;
  2. 裁剪会无条件删掉未消费的消息!XTRIM MAXLEN 1000 不管这些消息是否还在某个消费组的 PEL 里或者尚未投递。所以阈值必须远大于最大可能的 lag,并且要监控 XINFO GROUPS 的 lag。

Q9:Redis Stream 和 Kafka 的本质区别是什么?什么时候必须用 Kafka?

最本质的区别:存储介质。Stream 数据在内存,Kafka 在磁盘(顺序写)。

这一点决定了几乎所有其他差异:

  • 堆积成本:Stream 堆积 100GB 就要 100GB 内存(不可行);Kafka 堆积 100GB 只是磁盘(很便宜);
  • 保留时长:Stream 只能几小时到几天;Kafka 可以几周到几个月;
  • 扩展性:Stream 是单个 key,整个 stream 在一个节点上,无法像 Kafka 那样用 partition 并行。要提高吞吐只能自己分片(stream:orders:0/1/2),但会丢失全局顺序,且 rebalance 要自己实现。

其他关键差异:

  • 持久化保证:Kafka 的 acks=all + ISR 能保证不丢;Redis 的 AOF everysec 可能丢 1 秒,且异步主从复制在主节点宕机时会丢;
  • 语义:Stream 只有 At-Least-Once;Kafka 有事务 + 幂等生产者支持 Exactly-Once;
  • Rebalance:Kafka 自动;Stream 要靠 XAUTOCLAIM 手动实现;
  • 生态:Kafka 有 Connect、Streams、Flink 集成、Schema Registry。

必须用 Kafka 的场景:TB 级数据量、需要长期保存和任意重放、百万级 TPS、严格不丢消息、Exactly-Once、流处理生态。

Stream 的优势:延迟更低(微秒级)、运维成本低(已有 Redis 不用加组件)、实现简单。

Q10:用 zset 实现延迟队列有哪些坑?

坑一:必须用 Lua 保证"查询 + 删除"原子,否则多个消费者会重复消费:

local tasks = redis.call('ZRANGEBYSCORE', KEYS[1], 0, ARGV[1], 'LIMIT', 0, ARGV[2])
if #tasks > 0 then redis.call('ZREM', KEYS[1], unpack(tasks)) end
return tasks

(注意 unpack 有约 8000 个参数的栈限制,LIMIT 要控制在几百内。)

坑二:取出即删除,消费者崩溃则任务丢失。改进:租约机制——取出时不删除,而是把 score 改成"当前时间 + 处理超时",处理成功后再 ZREM 确认。这样超时未确认的任务会自动重新变成"到期"被别人取走。

坑三:member 唯一导致无法重复投递相同内容的任务——第二次 ZADD 相同 payload 只会更新 score。解决:给 payload 加唯一 ID。

坑四:轮询延迟 + 空轮询浪费。改进:用 ZRANGE key 0 0 WITHSCORES 看最近任务的到期时间,动态调整睡眠时长(但要能被新任务打断)。

坑五:所有任务在一个 zset 里形成大 key。改进:按时间分片(delay:queue:2026072913 按小时分,消费者同时轮询当前和上一小时)或按业务类型分 key。

Q11:notify-keyspace-events 的 K 和 E 有什么区别?

两种通知的频道和内容正好互换:

频道格式 消息内容
K(Keyspace) __keyspace@<db>__:<key> 事件名(如 set、expired)
E(Keyevent) __keyevent@<db>__:<event> key 名

执行 SET foo bar 且开启 KEA 时,会收到两条通知:

频道 __keyspace@0__:foo    内容 "set"
频道 __keyevent@0__:set    内容 "foo"

怎么选:

  • 关心某个特定 key 的所有变化 → 用 K:SUBSCRIBE __keyspace@0__:mykey;
  • 关心某类事件发生在哪些 key 上 → 用 E:SUBSCRIBE __keyevent@0__:expired(最常见)。

必须至少包含 K 或 E 之一,否则不发任何通知。A 是 g$lshzxetd 的别名(不含 m key-miss 和 n new-key)。

Q12:7.0 的分片 pub/sub 解决了什么问题?

问题:Redis Cluster 中,为了保证"订阅者连在任意节点都能收到消息",7.0 之前的 PUBLISH 会广播到集群里所有节点(通过集群总线)。

后果:一条消息在 10 节点集群里产生 10 倍的内部网络流量。高频 pub/sub 场景下集群总线可能被打满,影响正常的 gossip 通信和数据请求。

7.0 的分片 pub/sub(SSUBSCRIBE/SPUBLISH/SUNSUBSCRIBE):

  • 分片频道名按 CRC16 算 slot(和普通 key 一样,也支持 hash tag);
  • 消息只在负责该 slot 的主节点及其从节点之间传播,不再全集群广播;
  • 代价:订阅者必须连到正确的节点,客户端要处理 MOVED 重定向(客户端库需要支持)。

配套的查询命令:PUBSUB SHARDCHANNELS、PUBSUB SHARDNUMSUB。


小结

  • Pub/Sub 是"即时无状态转发":PUBLISH 推给当前在线的订阅者后就丢弃消息。不持久化、无 ACK、无消费组、堆积会断连、PUBLISH 占主线程、集群下 7.0 前全量广播。
  • Pub/Sub 只能用于允许丢消息的场景:配置热更新、缓存失效广播、监控推送、在线状态、锁释放通知(有超时兜底)。7.0 的 SPUBLISH/SSUBSCRIBE 分片 pub/sub 解决了集群广播放大问题。
  • RESP2 下订阅连接不能执行普通命令(客户端无法区分推送和回复),必须单独开连接;RESP3 的 > Push 类型解除了这个限制。
  • 键空间通知:K(频道带 key,内容是事件名)和 E(频道带事件名,内容是 key);notify-keyspace-events "Ex" 最常用。但不能用过期通知做延迟任务——触发时机不准(惰性+定期删除)、消息会丢(走 pub/sub)、广播导致重复消费。开 KEA 有 10%~30% 吞吐损失。
  • list 做队列的演进:RPOP 轮询(延迟/浪费)→ BRPOP(会丢消息)→ BRPOPLPUSH 可靠队列(但 list 元素没时间信息,看门狗无法判断超时;LREM 是 O(N))。固有缺陷:无消费组、无回溯、无进度观测、独占连接、无延迟无优先级。
  • Stream 是做 MQ 的唯一正确选择(5.0+)。生产者必须带 MAXLEN ~ 裁剪(用近似裁剪,且阈值远大于 lag);消费者三件事:启动时用 STREAMS key 0 恢复自己的 PEL、正常消费用 > + BLOCK、巡检 XAUTOCLAIM 认领超时消息 + 死信处理(RetryCount > 阈值则转死信并 ACK)。
  • 消费者名字必须稳定唯一(用 pod 名而非 UUID),否则 PEL 变孤儿、僵尸消费者堆积。XAUTOCLAIM 的 min-idle-time 必须 > 最大处理时间 × 2~3,否则造成重复消费。
  • 监控 XINFO GROUPS 的 lag(消费滞后)和 pending(未确认数)、XPENDING 的投递次数、消费者 idle。
  • XDEL 不释放内存(只打标记),XTRIM 会无条件删除未消费的消息。
  • Stream vs Kafka 的本质区别是存储介质(内存 vs 磁盘),由此决定了堆积成本、保留时长、扩展性(Stream 是单 key 在单节点,无 partition 并行)。Stream 是"够用的轻量 MQ",TB 级数据量/长期重放/百万 TPS/Exactly-Once/严格不丢 必须上 Kafka。
  • 延迟队列用 zset + Lua + 租约 + 时间分片;高可靠场景要用"Redis 快路径 + 数据库定时任务兜底"的双保险。