上一篇用租约锁协调“谁执行”,但锁不会保存任务、确认结果或重放失败。本篇从允许丢失和允许重复的代价出发比较 List、Pub/Sub 与 Stream:它们分别面向简单工作分发、在线广播和可恢复日志,API 看起来相似,离线补收、失败重投与积压观测能力却完全不同。
一、选型:先定义消息丢失与重复的代价
List 配合LPUSH/BRPOP能构成简洁 FIFO 队列,但消费者弹出后崩溃,任务已经消失。可用BLMOVE原子把任务移到 processing List,成功后移除,超时任务再回收;这需要自行维护确认、重试次数和死信。Pub/Sub 只把消息推给当时在线的订阅者,不持久化、不确认,断线期间消息直接错过,适合缓存失效通知和在线状态,不适合订单任务。
Stream 将条目保存在有序日志中,XADD产生递增 ID。消费组让多消费者分担消息,已投递未确认条目进入 Pending Entries List;成功后XACK,宕机任务可由其他消费者检查并认领。它提供至少一次处理基础,不等于业务恰好一次:消费者可能处理成功但在 ACK 前崩溃,消息会再次出现,所以副作用必须幂等。
二、原理:积压、背压与保留是一个系统
生产速度长期大于消费速度,任何持久队列最终都会耗尽内存。队列必须有最大积压、生产者限流、消费者扩缩容、消息保留和死信策略。Stream 的近似裁剪MAXLEN ~更高效,但可能略超目标长度;若消费组仍需读取被裁剪条目,盲目裁剪会留下无法恢复的处理状态。容量按每条序列化字节、峰值速率和最长恢复时间估算。
消息体应包含事件 ID、类型、模式版本、发生时间与业务引用,避免塞入巨大完整对象。消费者以事件 ID 或业务唯一键做幂等:先在最终数据库写入去重记录,再执行同一事务内的业务变更。只在 Redis 放一个“已处理”短 TTL key,过期后旧消息仍可能重复,而且 Redis 与数据库之间没有原子性。
三、实现:用状态机理解至少一次
下面程序模拟消息第一次处理完成但 ACK 丢失,随后被重新投递。幂等集合确保副作用只发生一次,而投递次数仍是两次。
fromcollectionsimportdeque stream=deque([{"id":"1710000000-0","order":"A100","amount":99}])pending={}processed=set()ledger=[]deliveries=0defhandle(message):globaldeliveries deliveries+=1event_id=message["id"]ifevent_idinprocessed:return"duplicate"ledger.append((message["order"],message["amount"]))processed.add(event_id)return"applied"message=stream.popleft()pending[message["id"]]=message first=handle(message)ack_was_lost=Trueassertack_was_lost retry=pending[message["id"]]second=handle(retry)pending.pop(message["id"])print(f"first={first}")print(f"second={second}")print(f"deliveries={deliveries}")print(f"side_effects={len(ledger)}")print(f"pending={len(pending)}")运行输出:
first=applied second=duplicate deliveries=2 side_effects=1 pending=0下面脚本建立 Stream、消费组,读取并确认一条消息。MKSTREAM允许空流建组;生产环境不要每次删除流,组创建应在部署迁移中幂等执行。
#!/usr/bin/env bashset-euopipefailredis_url="${REDIS_URL:-redis://127.0.0.1:6379/0}"stream='demo:orders'group='billing'consumer='worker-1'redis-cli-u"$redis_url"DEL"$stream">/dev/null redis-cli-u"$redis_url"XGROUP CREATE"$stream""$group"0MKSTREAM>/dev/nullmessage_id="$(redis-cli-u"$redis_url"--rawXADD"$stream"'*'event_id evt-001 order A100 amount99)"mapfile-treply<<(redis-cli-u"$redis_url"--rawXREADGROUP GROUP"$group""$consumer"COUNT1STREAMS"$stream"'>')printf'created_id=%s\n'"$message_id"printf'delivered_id=%s\n'"${reply[1]}"pending_before="$(redis-cli-u"$redis_url"--rawXPENDING"$stream""$group"|head-n1)"acked="$(redis-cli-u"$redis_url"--rawXACK"$stream""$group""$message_id")"pending_after="$(redis-cli-u"$redis_url"--rawXPENDING"$stream""$group"|head-n1)"printf'pending_before=%s\n'"$pending_before"printf'acked=%s\n'"$acked"printf'pending_after=%s\n'"$pending_after"redis-cli-u"$redis_url"DEL"$stream">/dev/null四、运维:恢复消费者而不是只重启进程
消费循环要区分新消息>与自己的 pending 历史。启动后先恢复可重试 pending,再阻塞读取新消息;长期无人处理的条目可通过XAUTOCLAIM转给健康消费者。认领阈值必须大于正常处理时长,否则两个消费者会同时处理慢任务。每次重试记录次数和最后错误,超过上限写入独立死信 Stream,并告警,而不是无限毒化主队列。
监控包括 Stream 长度、组 lag、pending 数量、最老 pending 空闲时间、生产/确认速率、重试和死信。仅看 Stream 长度会误判:保留窗口可能让已消费条目仍存在。消费者名称应稳定且可追踪,实例下线后清理无用消费者元数据前先确认其 pending 已转移。
Pub/Sub 订阅者要实现重连,并接受期间丢消息。若失效通知丢失会造成严重陈旧,可在消息中带版本,同时让本地缓存短 TTL 自愈。Redis 7 的分片 Pub/Sub 能限制 Cluster 广播范围,但仍不提供持久化语义。
五、验证:把重复、乱序和毒消息作为正常输入
集成测试应在业务提交后、ACK 前强杀消费者,确认消息重投且最终状态只变一次;让处理时间超过认领阈值,检查是否错误并发;注入无法反序列化的消息,验证会进入死信而非阻塞分区。多个 Stream 或多个生产者之间不要假设全局顺序,业务若要求同一订单顺序,应以聚合 ID 分流或用版本拒绝倒序事件。
性能测试必须带真实消息大小和消费者副作用。批量读取能提高吞吐,却延长单条确认时间并扩大失败重放范围。BLOCK时间要短于应用优雅停机预算,使消费者能检查取消信号。部署新模式时先发布兼容读取者,再发布新生产者。
本篇建立了可恢复的消息处理闭环。下一篇复用 Sorted Set、String 与原子脚本,实现排行榜、唯一访问计数和有窗口的限流计数,并处理并列与精度问题。
消息模式演进也需要发布顺序。消费者应先支持旧版和新版,再让生产者发送新版;字段新增提供默认值,字段删除至少跨过最大保留窗口。事件类型与模式版本分开记录,解析失败时保留原始消息摘要和错误原因,避免值班人员只能看到“消费失败”。对个人信息设置最短必要保留期,死信同样受数据合规约束,不能因为排障方便永久保存。
队列容量验收可以用公式反推:峰值每秒消息数乘单条字节数,再乘最长允许恢复秒数,得到最低积压字节;随后加入 Stream 元数据、pending 和安全余量。消费者扩容前确认瓶颈不在数据库,否则增加 worker 只会把下游压垮。降级时可拒绝低优先级生产、合并可折叠事件,关键订单消息则保持接收并触发更高等级告警。
参考来源
- Redis 官方文档:Streams
- Redis 官方文档:Pub/Sub
- Redis 官方文档:XAUTOCLAIM
👍 觉得有用就点个赞 + 收藏,方便回头查阅;有疑问直接在评论区留言,我看到都会回。
🚀 本文属于《Redis 应用实战》系列,持续更新,关注不迷路。
📌 文章里的代码都能直接跑。想要可直接 clone 的完整工程 + 配套部署脚本 / 踩坑清单?评论一声或发邮件到cj2664@qq.com,我免费发你。
如果你正好在做类似系统、或有工程化难题想找人做,也欢迎邮件聊一句——我按实际情况评估,能落地的就接单或出方案。评论和邮件都能直接找到我,不用跳别的平台。