Redis 应用实战(5):消息队列与发布订阅
上一篇用租约锁协调“谁执行”但锁不会保存任务、确认结果或重放失败。本篇从允许丢失和允许重复的代价出发比较 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 streamdeque([{id:1710000000-0,order:A100,amount:99}])pending{}processedset()ledger[]deliveries0defhandle(message):globaldeliveries deliveries1event_idmessage[id]ifevent_idinprocessed:returnduplicateledger.append((message[order],message[amount]))processed.add(event_id)returnappliedmessagestream.popleft()pending[message[id]]message firsthandle(message)ack_was_lostTrueassertack_was_lost retrypending[message[id]]secondhandle(retry)pending.pop(message[id])print(ffirst{first})print(fsecond{second})print(fdeliveries{deliveries})print(fside_effects{len(ledger)})print(fpending{len(pending)})运行输出firstapplied secondduplicate deliveries2 side_effects1 pending0下面脚本建立 Stream、消费组读取并确认一条消息。MKSTREAM允许空流建组生产环境不要每次删除流组创建应在部署迁移中幂等执行。#!/usr/bin/env bashset-euopipefailredis_url${REDIS_URL:-redis://127.0.0.1:6379/0}streamdemo:ordersgroupbillingconsumerworker-1redis-cli-u$redis_urlDEL$stream/dev/null redis-cli-u$redis_urlXGROUP CREATE$stream$group0MKSTREAM/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$consumerCOUNT1STREAMS$stream)printfcreated_id%s\n$message_idprintfdelivered_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)printfpending_before%s\n$pending_beforeprintfacked%s\n$ackedprintfpending_after%s\n$pending_afterredis-cli-u$redis_urlDEL$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 官方文档StreamsRedis 官方文档Pub/SubRedis 官方文档XAUTOCLAIM 觉得有用就点个赞 收藏方便回头查阅有疑问直接在评论区留言我看到都会回。 本文属于《Redis 应用实战》系列持续更新关注不迷路。 文章里的代码都能直接跑。想要可直接 clone 的完整工程 配套部署脚本 / 踩坑清单评论一声或发邮件到cj2664qq.com我免费发你。如果你正好在做类似系统、或有工程化难题想找人做也欢迎邮件聊一句——我按实际情况评估能落地的就接单或出方案。评论和邮件都能直接找到我不用跳别的平台。