返回归档
📮Redis

Redis Pub/Sub vs Redis Streams vs Asynq

一张表对照 Pub/Sub、Streams、Asynq 的投递语义、持久化、ACK、重试与回放。同一封 mail:send,三种失败完全不同。

文章目录

一句话主线:Pub/Sub、Streams、Asynq 不是三个「差不多的队列」——一个是在线广播,一个是可回放的事件日志,一个是带租约的任务状态机;选型先看投递合同,再看要不要延迟、去重、回放。

PUBLISH mail:send {"order":"42"}

Worker 正在重启。命令返回 0。邮件没了。

0 是当时收到消息的客户端个数,不是稍后投递的队列长度。同一封 mail:send order=42,换三条路,失败形态完全不同。本文目标只有一个:把这三个组件摆在同一张对比表上,说清各自合同,再决定用谁。

架构一眼看清

Pub/Sub 火忘即焚;Streams 日志加消费组;Asynq 任务状态机

  • Pub/Sub:发布者推到频道,在线订阅者立刻收到。没人听 → 消息直接消失。
  • Streams:消息先追加到持久化日志,再经消费组分发;要 XACK 才算完成。崩溃后靠 PEL 恢复。
  • Asynq:任务进多状态队列(Pending → Active → Retry / Archived)。Worker 用租约认领,崩溃后自动恢复。底层是 List + ZSET + Hash;取消偶借用 Pub/Sub;不用 Streams

三条路可以跑在同一个 Redis 上,但互不补漏:Pub/Sub 丢掉的那封,不会自动出现在 Stream 或 Asynq pending 里。

对比表

对比维度 Redis Pub/Sub Redis Streams Asynq
投递语义 At-most-once(最多一次) At-least-once(至少一次,配合 XACK At-least-once(至少一次 + 内置重试)
消息持久化 不持久化 持久化(可用 XTRIM 控制保留) 持久化(依赖 Redis AOF/RDB)
离线消息 直接丢失 可稍后消费 可自动恢复
消费者模型 广播(1:N,在线订阅者都收到) 广播,或消费组(组内负载均衡) 工作队列(一个任务只被一个 Worker 处理)
ACK 确认 有(XACK 隐式(成功处理后删除任务 key)
失败重试 需手动 XCLAIM / XAUTOCLAIM 内置自动重试 + MaxRetry + 退避
崩溃恢复 有(PEL + XCLAIM 有(Lease 过期后 Recoverer 捞回)
历史回放 支持(按消息 ID) 有限(Archived / Completed 可查看)
延迟 / 定时 需自己实现 原生(ProcessIn / ProcessAt
优先级队列 需多 Stream 模拟 原生(加权 / 严格优先级)
唯一性 / 去重 需自己实现 原生(Unique,窗口内不去重入队)
监控 基础(PUBSUB 等命令) 中等(XINFO 强(asynqmon、CLI、Prometheus)
使用复杂度 极低 中等(要自己管 ACK / claim) 对业务方低(库已封装状态机)
语言支持 全语言 全语言 主战场是 Go
底层 原生 Pub/Sub 原生 Stream List + ZSET + Hash 自建状态机

表里的「有 / 无」是组件合同,不是「能不能在外面再包一层」。用 Streams 也能做出延迟和去重,只是那一层要你自己写;Asynq 已经写进库里。

三个组件各自在干什么

Redis Pub/Sub:在线广播

官方 写明 at-most-once。发出去就结束。

PUBLISH 返回 0,邮件不落盘也不重送

  • 合同
    • at-most-once:断线、处理失败、当时没人听 → 都不会再送
    • PUBLISH 返回值 = 当时收到的客户端数(不是队列积压)
    • 无 backlog、无 ACK
  • 会怎样
    • Worker 重启窗口里返回 0 → 开篇那封邮件直接没了
  • 适合 / 不适合
    • 适合:缓存失效、在线状态(丢了可重建)
    • 不适合:mail:send(用户不会自己重试确认邮件)
    • 不进 keyspace:mail:send 不加环境前缀时,staging 也会吃到生产事件

Redis Streams:可回放的事件日志

Streams 是 append-only 日志。文档:先追加,再按消费组领取与确认。

XADD mail:events * order 42 type mail:send
XREADGROUP GROUP mailers w1 COUNT 1 STREAMS mail:events >
# ... 处理 ...
XACK mail:events mailers <entry-id>
  • 合同
    • at-least-once(配合 XACK
    • 可能被领两次 → 处理要幂等
  • 主路径
    • XADD → 落日志(单调 ID)
    • XREADGROUP → 进 PEL(Pending Entries List)
    • XACK → 离开 PEL
  • 崩溃时
    • 崩在 ACK 前:条目仍在 PEL
    • XCLAIM / XAUTOCLAIM 按空闲时间转交
    • 按消息 ID 可回放历史
  • 和「领任务」的区别
    • XREAD(不进组):无 ACK / 无 PEL → 回放或旁路读,不是领任务
    • 独立消费组:可各自读同一条 Stream(邮件组 vs 审计组)
    • 组内多 consumer:组内负载均衡
    • 忘记 XACK → PEL 只增不减、内存涨(不是 Redis 的 bug)
    • 延迟 / 去重 / 退避要自己建;别当成开箱任务队列

Asynq:带租约的任务状态机

Asynq 是 Go 侧常用的 Redis 任务队列。

  • 底层
    • List + ZSET + Hash 自建状态机
    • 不用 StreamsDiscussion #418:List / sorted set 够用)
    • 取消偶借用 Pub/Sub(asynq:cancel
  • 主路径
    • 入队:HSET 任务 Hash + LPUSH pending
    • 认领:RPOPLPUSH pending → active,并写入 lease ZSET(LeaseDuration = 30s
    • 成功:删任务 Hash(ACK 隐式)
  • 租约怎么理解
    • 30s 是心跳周期,不是「任务最长只能跑 30 秒」
    • heartbeater 会 ExtendLease
    • 心跳断了 → Recoverer 当崩溃处理
  • 崩溃路径
    • lease 过期 → retry ZSET
    • ForwardIfReady 到期后再 LPUSH 回 pending
  • 合同与能力
    • at-least-once(README:至少执行一次)
    • 删 key 前崩溃会再跑 → 按 order=42 做业务幂等
    • 超过 MaxRetry → archived
    • 原生:ProcessIn / ProcessAt、优先级、Unique、退避、asynqmon
    • asynq:cancel 只通知正在跑的 Worker
    • 还在 pending 的任务要出队删 Hash,不是发取消

同一封邮件,三种失败

路径 Worker 重启时 order=42 谁捞回来
Pub/Sub PUBLISH 返回 0,没了 没人
Streams 在日志 / PEL 里;忘了 XACK 则 PEL 涨 你自己 claim
Asynq lease 过期 → retry → 回 pending Recoverer + ForwardIfReady

逻辑拆开看:

  • 会不会丢
    • Pub/Sub:会(没人听 / 重启窗口)
    • Streams / Asynq:默认不丢(前提是真 ACK / 租约恢复在跑)
  • 谁负责捞
    • Pub/Sub:没人
    • Streams:你自己 XCLAIM / XAUTOCLAIM
    • Asynq:库内 Recoverer + ForwardIfReady
  • 会不会重复
    • Streams / Asynq:都会(at-least-once)→ 业务幂等都要写

怎么选

  • 选 Pub/Sub,当同时满足:
    • 只要实时广播
    • 丢了可以重建(缓存失效、在线提示)
  • 选 Streams,当同时满足:
    • 要可靠事件流
    • 要按 ID 回放,或要多服务独立消费同一条日志
    • 你愿意自己管 XACK / claim,并写幂等
  • 选 Asynq,当同时满足:
    • 要后台任务(发邮件、重试、延迟、去重)
    • 尤其是 Go 项目
    • 仍要写业务幂等(租约过期会再跑)

常见误判:

  • 把 Streams 当开箱任务队列
    • 延迟、退避、Unique、租约心跳、监控 UI → Streams 都要自建
    • Asynq 已经用另一套结构写过一遍
  • 用 Asynq 冒充事件总线
    • 「同一条日志被多组独立消费并回放」→ 工作队列模型对不上
    • 那是 Streams 的主场

收束:

  • 丢得起 → 广播(Pub/Sub)
  • 要历史 → 日志(Streams)
  • 要任务状态机 → Asynq(或同类库),别从 XADD 从零攒一套