用户支付成功后,订单接口还要做两件事:给用户发送通知,把“订单已支付”写进审计记录。最直觉的写法是全部同步完成,再向客户端返回成功。可通知供应商一抖,订单接口也跟着变慢;审计库短暂不可用,用户甚至会看到支付失败,刷新后却发现订单已经付过了。
于是我们把通知和审计改成异步任务:订单接口只提交订单,再发出一条 OrderPaid 事件,由后台消费者慢慢处理。请求确实变短了,但系统从此多出了一份必须说清楚的契约:事件什么时候算发出,消费者什么时候算处理完成,进程在两者之间崩溃会发生什么?
这正是本章的主线。我们不会把 Redis Stream 当成一串命令来背,也不会把它包装成“小号 Kafka”。我们会沿着订单 o1001 的一次支付,逐段拆开 List 的丢失窗口、Pub/Sub 的在线广播、Stream 的消费者组与待确认列表,再把重复投递、毒消息、积压和裁剪放回一个能上线的处理流程里。
异步只把工作移出了请求链路,没有把可靠性问题变没。同步调用的失败通常立刻暴露;异步系统的失败可能先沉到积压、待确认列表或死信里,几分钟后才被监控发现。设计队列时,先写故障语义,再写成功路径。
“发一条消息”听起来是同一件事,实际上至少有三种完全不同的关系。

如果只看“都能传字符串”,三种方案似乎可以互换;只要加入一次消费者宕机,差别就会立刻显现。下面的实验把同一条订单消息依次放进三种模型。切换模型后逐步执行“生产、领取、消费者崩溃、恢复”,观察消息到底还在不在,以及 Redis 是否知道谁领过它。
生产者可以 LPUSH queue:notify payload,消费者用 BRPOP queue:notify 5 阻塞等待。Redis 返回元素时,这个元素已经从 List 删除。假设 worker 刚收到 o1001,还没调用通知服务就退出,队列里没有消息,Redis 也没有一张“谁拿走了但没完成”的清单。这个从取出成功到业务完成之间的间隔,就是 List 最容易被忽略的丢失窗口。
List 并非只能做不可靠队列。可以用 BLMOVE 把元素从 queue:{notify}:ready 原子移动到 queue:{notify}:processing,成功后再用 LREM 删除,并由巡检程序把停留太久的任务移回 ready。问题是任务所有者、领取时间、投递次数、接管规则和死信都要由应用自己维护。需求长到这一步时,通常已经在手工重做 Stream 的一部分能力。
PUBLISH order:paid payload 会把消息推给当时已经订阅的连接。订阅者离线、网络断开或回调执行失败,消息不会留在频道里等待重发。它是“至多一次”的在线通知,不是待办箱。
这不代表 Pub/Sub 不可靠就没用。缓存失效提示、在线状态变化、开发环境日志尾随,都可能只关心此刻在线的接收者。真正危险的是把审计、结算、履约这类不能错过的工作误放进 Pub/Sub。ioredis 中订阅连接还会进入专门的订阅模式,工程上应使用 redis.duplicate() 为订阅单独建连接,不要把普通业务命令和订阅生命周期揉在一起。
Stream 是按 ID 递增的追加结构。XADD 写入后,消息仍留在 Stream;XREAD 可以从指定 ID 往后读;XREADGROUP 则让一个消费者组内的多个 worker 分摊新消息。消息被交给组内消费者后会进入 PEL,只有 XACK 才会移除这条组级待确认记录。
注意两件常被混在一起的事:XACK 不会删除 Stream 中的消息本体,它只说明某个消费者组处理完了;消息保留多久由 XTRIM、XADD MAXLEN/MINID 或删除策略决定。一个 Stream 也可以挂多个消费者组,例如 notification 与 audit 各自维护进度,所以两边快慢互不覆盖。
先更新订单数据库,再 XADD,会有一个很现实的故障:数据库提交成功后进程退出,事件没有发出。反过来先 XADD 再提交订单,消费者可能先看到一个最终回滚的订单。Stream 能管理已经进入 Redis 的消息,却不能自动把关系型数据库提交和 Redis 写入合成一个事务。
对“订单已支付”这种不能悄悄漏掉的事件,更稳妥的做法是事务型 Outbox。订单表和 outbox 表在同一个数据库事务里提交,由独立 relay 扫描未发布记录并写入 Stream。
支付请求在数据库事务中更新 orders.status = 'PAID',同时向 outbox_events 插入事件。两条 SQL 要么一起提交,要么一起回滚。
relay 读取未发布事件,用 event_id、order_id、event_type 和版本号调用 XADD。事件负载只放消费者真正需要的稳定字段,不把整张订单对象随意序列化进去。
XADD 成功后,relay 再把 outbox 行标记为已发布。如果它在两步之间退出,同一个 event_id 可能再次写进 Stream,所以消费者仍要幂等。
BEGIN;
UPDATE orders
SET status = 'PAID', paid_at = NOW()
WHERE id = 'o1001' AND status = 'PENDING';
INSERT INTO outbox_events(event_id, aggregate_id, event_type, payload)
VALUES ('evt_o1001_paid_v1', 'o1001', 'OrderPaid', '{"amount":"99.00"}');
COMMIT;事件名表达已经发生的事实,所以用 OrderPaid,而不是含糊的 handleOrder。event_id 是投递去重键,order_id 是业务聚合标识,event_version 用来处理以后字段升级。Stream 自动生成的 ID 负责流内位置,不应兼任跨系统业务幂等键。
如果业务允许偶尔漏一条“营销提醒”,数据库提交后直接 XADD 也可能是合理取舍;如果事件驱动的是审计、履约或账务衔接,就应使用 Outbox、数据库变更订阅或另一种明确的可靠发布机制。选择可以简单,但失败代价必须写在设计里。
我们沿用课程中的订单流:stream:{orders}:events 保存订单事件,fulfillment 是履约消费者组,组内有 worker-a 与 worker-b。花括号不是装饰;若以后在 Redis Cluster 中用脚本同时操作订单流与死信流,可以让相关 key 使用同一个哈希标签。
第一次创建组时,起始 ID 决定要不要处理已有历史。0-0 表示从最早保留记录开始,$ 表示从建组时的末尾开始,只处理之后新增的消息。创建动作应由部署初始化或幂等启动代码完成,并处理组已存在的情况。
import Redis from "ioredis";
const options = {
host: "127.0.0.1",
port: 6379,
maxRetriesPerRequest: 1,
};
const commands = new Redis(options);
const reader = commands.duplicate(); // BLOCK 读取使用专用连接
> 只在消费者组读取里有特殊含义:取尚未投递给本组任何消费者的新消息。它不是“最后一条 Stream ID”的通用写法。若 worker 要查看分配给自己的未确认历史,应传一个具体 ID,例如 0,而不是继续传 >。
Stream ID 通常形如 1753212345678-0,前半段近似毫秒时间,后半段用于同一毫秒内排序。它提供稳定的流内次序与读取位置,但不等于业务发生时间:网络延迟、补发和不同生产者的时钟都会让两者不完全一致。需要按业务时间判断时,应在字段里另外保存 occurredAt。
阻塞读取不要占用普通命令连接。一个正在执行 XREADGROUP ... BLOCK 的连接要等消息或超时后才能返回,同一连接后面的 GET、XACK 也会排队。ioredis 可以用 duplicate() 建专用 reader,并把 BLOCK 设成有限值,让进程定期检查退出信号、连接状态和配置变化。
同一个 Stream 中,消息 ID 按追加顺序递增;消费者组把不同消息交给多个 worker 后,业务完成顺序却可能变化。worker-a 先拿到“订单已支付”,处理通知用了 8 秒;worker-b 后拿到“订单已退款”,1 秒就处理完,用户可能先看到退款通知,再看到支付通知。消费者组保证的是分摊与待确认管理,不是多个 worker 之间的串行完成。
若同一订单的事件必须按版本执行,可以在消息里携带 aggregateId 与 aggregateVersion,消费者只接受期望的下一个版本,发现缺口时暂缓并告警。吞吐更大时,可以按 orderId 稳定分片到多个 Stream key,同一订单始终进入同一分片,并让每个分片保持受控的串行度。这已经接近 Kafka 的按 key 分区思路,只是分片映射、扩缩容和重平衡都要由应用承担。不要用一个覆盖全局的分布式锁强行恢复顺序,那会把并行消费重新压成单点。
停机也属于顺序的一部分。进程收到退出信号后,应先停止领取新消息,再等待正在执行的数据库事务结束,成功的任务做 XACK,超过关停预算的任务保持 pending 交给后续接管。直接杀掉进程不会破坏 Stream 本体,但会制造大量需要等待 idle 阈值的 PEL;若每次发布都这样做,接管告警会变成常态噪音。
课程项目中,同一组一次读入两条固定 ID 的订单事件,确认前后得到下面的结果:
读取消息数: 2
确认前 pending: 2
XACK 成功数: 2
确认后 pending: 0这组数字只证明 PEL 的变化,不证明通知已经送达用户。Redis 知道的是“是否收到确认”,业务系统仍要自己定义“什么条件下允许确认”。

worker-a 用 XREADGROUP 读到消息后,Redis 会在组的 Pending Entries List 中记录消息 ID、当前所有者、空闲时长和投递次数。worker-a 正常处理并 XACK,记录离开 PEL;如果进程退出,记录会一直留着,除非被确认、被接管,或运维策略另行清理。
XPENDING stream group 适合看摘要:待确认总数、最小与最大 ID、各消费者持有数量。扩展形式可以按范围取明细,看到 owner、idle 和 delivery count。巡检程序发现某条消息空闲超过阈值后,可用 XAUTOCLAIM 把它转给 worker-b。这个阈值不能拍脑袋设成 5 秒:如果正常生成发票的 P99 是 12 秒,5 秒接管会让两个 worker 同时处理同一订单。
XPENDING stream:{orders}:events fulfillment
XAUTOCLAIM stream:{orders}:events fulfillment worker-b 30000 0-0 COUNT 50XAUTOCLAIM 像扫描游标一样返回下一次起点,调用方要持续推进游标,直到返回 0-0。它只认领 idle 足够久的 pending 消息,并会增加投递次数;它不会自动决定重试上限,也不会自动把失败消息送进死信。接管还会重置 idle,因此恢复程序必须避免多个实例毫无协调地高频全表扫描。
下面的实验可以创建事件、让 worker-a 在确认前崩溃、推进 idle 时间,再由 worker-b 接管。继续制造失败,会看到投递次数变化;开启幂等后,即使消息重复到达,业务副作用也只保留一份。
消费者最安全的基本顺序是:读取消息 → 完成业务副作用 → XACK。若反过来先确认再发通知,确认后进程退出会永久漏通知。可按安全顺序执行时,另一个窗口又出现了:通知已经发送,进程在 XACK 前退出;worker-b 接管后会再次处理,于是用户收到两条“订单已支付”。
这不是 Stream 的 bug,而是至少一次投递的直接结果。Redis 无法把外部短信接口、审计数据库事务和 XACK 包进同一个原子提交。所谓“恰好一次”,最终必须落到业务副作用是否能识别同一个 eventId。

一种接地气的做法,是在业务数据库建 processed_events(event_id primary key),并让去重记录与审计记录、通知任务在同一数据库事务中提交。重复事件再次到达时,唯一约束让事务知道它已处理过,消费者跳过副作用,然后补做 XACK。
async function handleOrderPaid(messageId, fields) {
const event = Object.fromEntries(
Array.from({ length: fields.length / 2 }, (_, index) => [
fields[index *
这里没有在数据库事务中直接调用通知供应商,因为持有数据库事务等待外部网络,既扩大锁时间,也仍然得不到跨系统原子性。我们把通知变成有唯一 event_id 的本地任务,再由发送器调用支持幂等键的供应商;若供应商不支持幂等,就保存明确的发送状态并准备人工对账。幂等不是简单 SETNX 一次:去重记录若先写成功、业务随后失败,也会错误地吞掉重试,因此去重状态要与本地副作用一起提交。
异步之后,生产者与消费者不会同时发布。今天给 OrderPaid 增加 couponAmount,线上仍可能有旧 worker 在读取。事件字段应尽量向后兼容:新增字段提供默认含义,金额用明确单位或字符串避免浮点歧义,枚举扩展时让消费者能识别“未知值”,删除和改名则通过新版本过渡。消费者遇到不支持的 eventVersion,应该把原因写进错误上下文并转入可治理的重试或死信流程,不能把未知字段静默解释成零。
事件也不宜无限膨胀。只传 orderId 会迫使每个消费者回查数据库,放大耦合与读压力;把订单所有字段复制进去,又会泄露不必要信息并让版本演进困难。更稳妥的边界是:事件包含判断这个事实所需的稳定快照,例如订单号、用户标识、金额、币种、发生时间与版本;敏感地址、完整支付凭证等字段由有权限的服务按需读取。消息契约一旦被多个组使用,就应像接口一样做兼容测试。
不要用 Stream ID 代替业务事件 ID。Outbox relay 重发同一个业务事件时,可能产生两个不同的 Stream ID;如果消费者只按 Stream ID 去重,重复副作用仍会发生。稳定的 eventId 应由业务事件在第一次创建时生成,并穿过每一次转发。
通知服务超时、数据库连接池暂时耗尽,通常属于临时故障;手机号格式错误、事件缺少必填字段、消费者版本不认识 payload,则更像永久故障。把两类错误都立即重试,会让一条坏消息占满 worker,拖住后面的正常订单。
可以把处理策略拆成三层:第一次失败记录错误分类并保留 pending;临时故障按指数退避加随机抖动再次尝试;投递次数超过阈值或确认是永久错误后,将原消息写入 stream:{orders}:dead,带上原 Stream ID、eventId、失败原因、投递次数与失败时间,再确认原消息。死信操作还应产生告警和工单,而不是把消息换个 key 后继续无人负责。

XAUTOCLAIM 只负责所有权转移,并不提供延迟调度。若重试需要 10 秒、1 分钟、5 分钟的退避,可以用 Sorted Set 保存下一次执行时间,由调度器到点再写入重试 Stream;也可以选择已经内置延迟、退避和死信策略的任务框架。无论采用哪种方式,都要给同一 eventId 保留统一的累计尝试次数,不能每次换 Stream 后又从 1 开始。
把消息写入死信再 XACK 原消息是两步。若要求 Redis 内部不出现“死信已写但原消息未确认”的窗口,可以用事务或短 Lua 脚本把两条 Redis 命令放在一起;Cluster 下两个 key 必须落在同一槽,所以示例才使用 {orders} 哈希标签。但脚本仍然不能覆盖数据库或外部接口,消费者幂等依旧不可省。
死信至少需要一条完整回路:告警能找到负责团队;后台能查看脱敏后的 payload 与错误;修复数据后可以用同一个 eventId 重放;重放成功后记录操作者、时间和结果。没有这些动作,“死信队列”只是一个更隐蔽的消息墓地。
假设订单高峰每秒产生 2000 条事件,通知消费者只能处理 1200 条。接口仍然很快,Redis 却每秒多出 800 条积压。一个小时后就是 288 万条。Stream 的低延迟并不会改变这个算术;当生产速率长期高于消费能力,唯一结果是内存持续上涨。
阻塞与批量参数要一起调。BLOCK 5000 避免无消息时空转,COUNT 20 减少网络往返,但它不是服务端承诺每次一定给 20 条。批次越大,单个 worker 一次持有的 pending 越多,崩溃后的接管量也越大。消费者还要限制本地并发,不能一次读 1000 条后同时打向同一个审计数据库。
排查积压时,至少看下面这些信号:
Stream 默认会持续增长,必须显式决定保留边界。XADD ... MAXLEN ~ 100000 用近似长度裁剪,开销通常比精确裁剪更平滑;XTRIM ... MINID ~ <id> 更适合按时间位置保留。裁剪不是消费确认:若保留窗口短于最慢消费者的处理与回放窗口,消息本体可能已被删,PEL 中却还留着引用,接管时只剩 ID 没有 payload。
因此保留策略要基于“最大可接受积压时间 + 故障恢复时间 + 人工重放窗口”,并给容量留余量。较新的 Redis 版本提供了更细的消费者组引用处理选项,但升级命令语义之前仍要确认所有组的兼容性。最稳妥的原则没有变:不要用裁剪代替消费治理,也不要假设已 XACK 的消息会自动从 Stream 消失。
消息流最好不要与可随意淘汰的缓存混在同一个内存池里。若实例使用 allkeys-lru 一类策略,内存压力可能把 Stream key 当普通 key 淘汰。即便启用了 AOF、复制或高可用,也要明确可接受的数据丢失窗口;Redis 的故障恢复能力不等于订单数据库与事件流之间天然一致。
Stream 与 Kafka 都有追加记录、消费者组和历史读取,所以功能表很容易把它们画成“大同小异”。真正影响选型的,是系统预期承受多少数据、保留多久、怎样扩展并行度、如何重放,以及团队愿意运营什么。

所以“不能丢就一定 Kafka、允许丢才用 Stream”太粗糙。一个配置了持久化与副本、短积压、消费者幂等、有人值守的 Stream,可以承接相当严肃的内部任务;一个副本数为 1、生产者确认配置随意、没人监控的 Kafka,也不会自动变可靠。
更实用的分界是:如果服务已经依赖 Redis,消息规模有限,积压按分钟或小时计算,只需要几个消费者组,历史回放不是核心产品能力,Stream 往往够用。若事件要保留数天乃至更久,峰值吞吐与积压很大,需要按业务 key 分区扩展,大量团队会独立重放同一份日志,那么这些条件本身就在把系统推向 Kafka。不要等 Redis 内存报警后,才把“以后可能换”当迁移方案。
下面的选型台不会只吐出一个神秘分数。调整峰值写入、积压时长、保留天数、消费者组数量、按 key 顺序、历史重放与运维能力,它会逐项解释是哪一条需求在给哪种方案施压。
走到这里,可以把订单事件链路压缩成一条可执行的检查顺序:订单与 Outbox 同事务提交;relay 用稳定 eventId 写 Stream;每类下游使用独立消费者组;worker 以专用连接阻塞读取;本地副作用和去重记录同事务提交;最后 XACK;巡检用 XPENDING 与 XAUTOCLAIM 接管超时消息;重试有退避和上限;死信有人处理;保留策略大于恢复窗口;lag、PEL、最老消息年龄与内存都有告警。
任何一项写不清,都不要用“消息队列保证了可靠性”带过。队列能保存状态,却不知道业务完成的定义;消费者组能记录 pending,却不会替你识别毒消息;XACK 能结束一次投递,却无法撤销已经发出的短信。可靠性来自这一整条协议,而不是某个 Redis 命令。
下一章讨论事务与原子性时,我们会继续追问同一个问题:哪些步骤真的能被 Redis 合成一个不可穿插的操作,哪些跨系统结果仍然只能靠幂等、补偿和对账收口。
定时对账检查长时间未发布的 outbox、Stream 最新 ID 与消费者组积压。没有对账的 Outbox,只是把“可能漏消息”换成“可能永远卡在表里”。