目录

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:techPSUBSCRIBE 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 通用命令(DELEXPIRERENAME 等)
$ 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 的别名(不含 mn

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

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 的所有变化 → 用 KSUBSCRIBE __keyspace@0__:mykey
  • 关心某类事件发生在哪些 key 上 → 用 ESUBSCRIBE __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. 必须设置 MaxLenMinID 裁剪,否则内存会无限增长;
  2. Approx: true(对应 ~)而不是精确裁剪:近似裁剪只在能整个删掉一个 listpack 节点时才删(O(1) 摊还),精确裁剪要逐条删(可能 O(N) 阻塞);
  3. 裁剪阈值要留足余量:如果消费组有滞后(lag),MaxLen 太小会直接丢掉未消费的消息。要监控 XINFO GROUPSlag
  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-0worker-1),保证重启后名字不变。用 XGROUP DELCONSUMER 清理确定不再回来的消费者。

(3)XAUTOCLAIMMinIdle 怎么设

必须 > 单条消息的最大处理时间。如果消息处理要 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:0stream: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 里待了多久,看门狗进程没法判断超时;
    • LREMO(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-0worker-1)这种稳定标识。

Q7:XAUTOCLAIMmin-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 vMINID ~ <时间戳>
  2. ~(近似裁剪)而不是 =(精确裁剪):近似裁剪只在能整个删掉一个 listpack 节点时才删(O(1) 摊还),精确裁剪要逐条删可能 O(N) 阻塞主线程;
  3. 定期 XTRIMXTRIM s MINID ~ <7天前时间戳>

两个危险点

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

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> 事件名(如 setexpired
E(Keyevent) __keyevent@<db>__:<event> key 名

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

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

怎么选

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

必须至少包含 KE 之一,否则不发任何通知。Ag$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 SHARDCHANNELSPUBSUB 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 变孤儿、僵尸消费者堆积。XAUTOCLAIMmin-idle-time 必须 > 最大处理时间 × 2~3,否则造成重复消费。
  • 监控 XINFO GROUPSlag(消费滞后)和 pending(未确认数)XPENDING 的投递次数、消费者 idle
  • XDEL 不释放内存(只打标记),XTRIM 会无条件删除未消费的消息
  • Stream vs Kafka 的本质区别是存储介质(内存 vs 磁盘),由此决定了堆积成本、保留时长、扩展性(Stream 是单 key 在单节点,无 partition 并行)。Stream 是"够用的轻量 MQ",TB 级数据量/长期重放/百万 TPS/Exactly-Once/严格不丢 必须上 Kafka。
  • 延迟队列用 zset + Lua + 租约 + 时间分片;高可靠场景要用"Redis 快路径 + 数据库定时任务兜底"的双保险。