别再手动封装List当队列了!Redis Stream消息队列,原理+实战一篇讲透
别再手动封装List当队列了Redis Stream消息队列原理实战一篇讲透 前言学消息队列时你可能也经历过这个阶段——先是在项目里用BlockingQueueThreadPoolExecutor凑合着解耦结果应用一重启队列里的消息全没了。你开始想有没有一种队列能让多个JVM实例共享消息而且挂掉也不丢数据然后你搜到了 Redis Stream。再然后你反反复复看了好几篇文章XADD、XREAD、XGROUP、XREADGROUP、XACK、XPENDING……十几个命令长得差不多每个的角色是什么、什么时候该用哪个脑子全是浆糊。这篇文章帮你把这根线从头捋到尾。 环境说明组件版本/说明OSWindows / macOS / LinuxRedis6.2Stream 从 5.0 开始支持Java11客户端Jedis 5.x 或 LettuceIDE任意1. 什么是消息队列消息队列Message Queue是一种跨进程的通信方式。生产者把消息发到队列后就立即返回消费者可以是任何进程、任何语言、任何机器异步地从队列里取出消息并处理。生产者和消费者彼此不知道对方的存在它们只和队列打交道。核心价值两点异步削峰生产者不用等消费者处理完再继续干活瞬时流量先打到队列里缓存消费者按自己的节奏处理独立扩缩生产者和消费者的实例数可以单独调整流量高峰时单独加消费者即可2. 为什么需要消息队列JDK 内置的BlockingQueue在同一个 JVM 内确实能实现生产者-消费者模式但受限于进程边界。消息队列弥补了它的短板对比维度JDK 内置队列BlockingQueue消息队列Redis Stream / RabbitMQ 等进程隔离❌ 同一个 JVM 内✅ 任意进程任意语言持久化❌ 进程重启即丢失✅ RDB/AOF 持久化Redis可靠消费❌ 无 ACK 机制✅ ACK 重试削峰填谷❌ 内存有限容易 OOM✅ 队列长度可配置/可限流消费者扩展⚠️ 需自己写分配逻辑✅ Consumer Group 天然支持总结下来需要消息队列的理由有四条解耦— 生产者和消费者不必在同一进程、同一语言、同一机器上削峰填谷— 瞬时流量缓存到队列消费者按节奏处理后端不会被冲垮可靠投递— 持久化到磁盘 ACK 机制确保每条消息至少被处理一次可扩展— 消费者实例数可以随时增加生产者完全无感3. 为什么是 Redis StreamConsumer GroupRedis 里可以做队列的数据结构不止 Stream——还有 ListBLPOP和 Pub/Sub。Stream 的 Consumer Group 是唯一同时满足持久化 ACK 确认 消费组负载均衡的方案特性ListBLPOPPub/SubStreamConsumer Group持久化✅ RDB/AOF❌ 不持久化✅ RDB/AOF消息确认❌ POP即删除❌ 发完即忘✅ XACK 确认消息回溯❌ 消费完就没了❌ 不存储✅ 可从头重读多消费者负载均衡❌ 竞争模式谁抢到算谁的✅ 广播所有订阅者都收到✅ 组内每条消息只给一个人消费组❌ 需要自己实现❌ 无此概念✅ 原生支持Consumer Group 的三个核心概念概念作用last-delivered-id组级别的进度指针标记最新投递到哪条消息PELPending Entry List整个组的共享数据结构每条消息标记归属于哪个消费者Consumer 名称组内区分消费者用于故障转移时 XCLAIMPEL 不是每个消费者独享的——它是一个组级别的共享数据结构。组内所有未 ACK 的消息都存在这里只是每条消息会通过consumer_name字段标记当前由哪个消费者负责。XPENDING无过滤时返回整个组的待确认消息按消费者过滤只是在同一个结构上做筛选。4. Redis Stream Consumer Group 工作原理全局 PEL组级别共享数据结构正常流程Consumer A 崩溃生产者 XADD 消息消息存入 StreamConsumer Group 分配Consumer AConsumer BConsumer C消息 M1consumer: AConsumer A 处理并 XACK从全局 PEL 移除Consumer B XCLAIM 接管归属变为 B处理并 XACK4.1 创建组XGROUP CREATE mystream mygroup $ [MKSTREAM]$表示从最新的消息开始消费不处理历史消息0则表示从头开始消费4.2 生产消息XADD mystream * field1 value1 field2 value2*让 Redis 自动生成 ID格式为时间戳-序号如1721645023000-0消息一旦写入不可修改append-only 结构4.3 消费消息消费者通过XREADGROUP读取消息有两种模式模式参数含义读新消息只读取尚未投递给组内任何消费者的新消息读历史重试具体 ID读取自己 PEL 中已经投递但未 ACK 的消息Redis不是让消费者竞争消息——消息的分配由服务端决定。当消费者用参数调用XREADGROUP时Redis 把尚未投递的消息按轮询方式分配给当前请求的消费者。这更像是老师按顺序点名而不是学生抢答。4.4 确认处理XACK mystream mygroup 1721645023000-0确认后消息从全局 PEL 中移除。如果不 XACK这条消息会一直留在 PEL 里。4.5 故障恢复假设 Consumer A 拿到消息后崩溃了没来得及 XACK。这条消息会一直留在全局 PEL 中归属为 Consumer A。其他消费者可以用XPENDING发现滞留消息再用XCLAIM把消息的归属权转移给自己处理最后 XACKXPENDING mystream mygroup XCLAIM mystream mygroup consumerB 60000 1721645023000-0XCLAIM不是把消息从一个消费者的私有列表搬到另一个——它只是在全局 PEL 中把consumer_name字段从 A 改成 B。这解释了为什么 XCLAIM 轻量且数据不会丢失。这样就实现了至少一次投递的语义消息要么被 ACK要么超时后被其他消费者接管重试。5. Redis Stream 在 Java 中的实现以下是最小可运行的生产者-消费者完整实例基于 Jedis 5.x。5.1 依赖dependencygroupIdredis.clients/groupIdartifactIdjedis/artifactIdversion5.1.0/version/dependency5.2 生产者创建组并生产消息importredis.clients.jedis.*;importredis.clients.jedis.params.XAddParams;importjava.util.*;publicclassStreamProducer{publicstaticvoidmain(String[]args){StringstreamKeymystream;StringgroupNamemygroup;try(JedisjedisnewJedis(localhost,6379)){// XGROUP CREATE mystream mygroup $ MKSTREAMtry{jedis.xgroupCreate(streamKey,groupName,null,true);System.out.println(Group created: groupName);}catch(Exceptione){System.out.println(Group already exists);}// 生产 5 条消息for(inti1;i5;i){MapString,StringbodynewHashMap();body.put(id,String.valueOf(i));body.put(msg,Hello MQ #i);// XADD mystream * id 1 msg Hello MQ #1StringmessageIdjedis.xadd(streamKey,XAddParams.xAddParams().id(*),body);System.out.println(Produced: messageId - body);}}}}5.3 消费者Consumer Group 模式importredis.clients.jedis.*;importredis.clients.jedis.params.XReadGroupParams;importjava.util.*;publicclassStreamConsumer{publicstaticvoidmain(String[]args){StringstreamKeymystream;StringgroupNamemygroup;StringconsumerNameconsumer-A;try(JedisjedisnewJedis(localhost,6379)){// XREADGROUP GROUP mygroup consumerA BLOCK 2000 COUNT 10ListMap.EntryString,ListStreamEntryresultsjedis.xreadGroup(groupName,consumerName,XReadGroupParams.xReadGroupParams().block(2000).count(10),true,// true 用 模式读新消息newStreamEntry(streamKey,null));if(resultsnull||results.isEmpty()){System.out.println(No new messages.);return;}for(Map.EntryString,ListStreamEntrystream:results){Stringkeystream.getKey();for(StreamEntryentry:stream.getValue()){StringmsgIdentry.getID().toString();MapString,Stringdataentry.getFields();System.out.println(Received [key]: msgId - data);// XACK mystream mygroup messageIdlongackedjedis.xack(streamKey,groupName,entry.getID());System.out.println( ACKd: msgId (ackedacked));}}}}}5.4 故障恢复检查 PEL 并接管importredis.clients.jedis.*;importredis.clients.jedis.params.XPendingParams;importjava.util.*;publicclassPELInspector{publicstaticvoidmain(String[]args){StringstreamKeymystream;StringgroupNamemygroup;try(JedisjedisnewJedis(localhost,6379)){// XPENDING mystream mygroupStreamPendingSummarysummaryjedis.xpending(streamKey,groupName);System.out.println(Total pending: summary.getTotal());System.out.println(Earliest ID: summary.getMinMessageId());System.out.println(Latest ID: summary.getMaxMessageId());// XPENDING mystream mygroup - 10ListStreamPendingEntrypendingListjedis.xpending(streamKey,groupName,XPendingParams.xPendingParams().start(-).end().count(10));for(StreamPendingEntryentry:pendingList){System.out.printf(Consumer: %s, Pending: %d%n,entry.getConsumerName(),entry.getPending());}// XCLAIM: 将超时 60s的未 ACK 消息转移给自己StringmyConsumerconsumer-B-backup;for(StreamPendingEntryentry:pendingList){if(entry.getPending()0){ListStreamEntryclaimedjedis.xclaim(streamKey,groupName,myConsumer,60000,null,entry.getEarliestMessageId());System.out.println(Claimed: claimed);}}}}}Jedis API 速查Jedis 方法对应 Redis 命令作用jedis.xadd()XADD生产消息jedis.xgroupCreate()XGROUP CREATE创建消费组jedis.xreadGroup()XREADGROUP消费者读取消息jedis.xack()XACK确认消息已处理jedis.xpending()XPENDING查询待确认消息jedis.xclaim()XCLAIM故障转移接管未 ACK 消息6. 总结要点回顾消息队列的核心价值是解耦 削峰 可靠投递Redis Stream Consumer Group 相比 List/Pub/Sub 的不可替代优势持久化 ACK 确认 故障转移工作流程XADD 生产 - Consumer Group 自动分配 - 消费者处理 - XACK 确认 - 从 PEL 移除PEL 是组级共享的全局数据结构每条消息标记所属消费者宕机未 ACK - XPENDING 发现 - XCLAIM 转移归属是故障恢复的标准流程XCLAIM 不搬数据只改归属——在全局 PEL 中修改consumer_name字段轻量且可靠Java 中 Jedis 5.x 的 API 直接对应 Redis 命令熟悉命令名称即可上手延伸思考Consumer Group 的“至少一次投递”语义很强但代价是消费者需要处理消息重复的可能XCLAIM 后可能被两个消费者处理。如果你在写金融/支付类的逻辑建议在业务侧做幂等——比如在消息体里加一个唯一 ID消费前去数据库查重。 参考资料Redis Stream 官方文档Jedis GitHubRedis Stream 设计原理antirez 原版介绍© 本文为原创内容转载请注明出处。如果这篇文章对你有帮助欢迎点赞 、收藏 ⭐、关注 ➕你的支持是我持续输出的动力