一句话主线:Pub/Sub、Streams、Asynq 不是三个「差不多的队列」——一个是在线广播,一个是可回放的事件日志,一个是带租约的任务状态机;选型先看投递合同,再看要不要延迟、去重、回放。
PUBLISH mail:send {"order":"42"}
Worker 正在重启。命令返回 0。邮件没了。
0 是当时收到消息的客户端个数,不是稍后投递的队列长度。同一封 mail:send order=42,换三条路,失败形态完全不同。本文目标只有一个:把这三个组件摆在同一张对比表上,说清各自合同,再决定用谁。
架构一眼看清

- 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。发出去就结束。
- 合同
- at-most-once:断线、处理失败、当时没人听 → 都不会再送
PUBLISH返回值 = 当时收到的客户端数(不是队列积压)- 无 backlog、无 ACK
- 会怎样
- Worker 重启窗口里返回
0→ 开篇那封邮件直接没了
- Worker 重启窗口里返回
- 适合 / 不适合
- 适合:缓存失效、在线状态(丢了可重建)
- 不适合:
mail:send(用户不会自己重试确认邮件)
- 坑
- 不进 keyspace:
mail:send不加环境前缀时,staging 也会吃到生产事件
- 不进 keyspace:
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) - 可能被领两次 → 处理要幂等
- at-least-once(配合
- 主路径
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 自建状态机
- 不用 Streams(Discussion #418:List / sorted set 够用)
- 取消偶借用 Pub/Sub(
asynq:cancel)
- 主路径
- 入队:
HSET任务 Hash +LPUSHpending - 认领:
RPOPLPUSHpending → 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从零攒一套