java用redis stream实现轻量级消息队列,核心是利用其持久化、消费者组和手动ack机制,支持至少一次语义,适合日均百万级以下场景;生产用xadd写入有序结构化消息,消费需通过xreadgroup绑定消费者组,失败消息由xpending和xclaim保障可靠性,运维需控制stream长度并监控pending数量。

Java 用 Redis Stream 实现轻量级消息队列,核心在于利用其持久化、消费者组和手动 ACK 机制,避免传统 List 或 Pub/Sub 的丢消息、无状态、难追溯等问题。它不是“模拟队列”,而是原生支持消息可靠投递的结构,适合日均百万级以下、要求“至少一次”语义的场景。
生产消息:用 XADD 写入带唯一 ID 的结构化数据
Stream 消息本质是键值对组成的日志条目,每条自动或手动分配唯一 ID(如 1640995200000-0),天然有序、可回溯。
- 推荐使用 Spring Data Redis 的
StreamOperations,避免手拼命令 - 消息体建议用 Map
,字段名语义清晰(如 "order_id","status") - ID 用
*让 Redis 自动生成,保证时序与唯一性;特殊场景(如幂等重发)可指定 ID - 示例:
Map<string string> message = Map.of("order_id", "ORD-20260728-001", "event", "paid"); redisTemplate.opsForStream().add(StreamRecords.string(message).withStreamKey("order_events"));</string>
消费消息:用 XREADGROUP 绑定消费者组与具体消费者
必须通过消费者组(Consumer Group)消费,才能启用 ACK 和 Pending 列表——这是可靠性的基石。
Java JDK 25 来自 OpenJDK 官方归档,版本为 JDK 25,本条下载地址已指向官方 Windows x64 zip 安装包直链,适合调试旧项目或兼容旧版 Java 运行环境。
- 首次创建组需显式调用
XGROUP CREATE order_events pay_group $ MKSTREAM($表示只消费后续消息) - 消费时指定组名 + 消费者名(如
"worker-1"),Redis 自动做组内负载均衡 - 用
count控制单次拉取条数(建议 10–50),用block实现阻塞等待,避免空轮询 - 示例(Lettuce 客户端):
ReadOffset offset = ReadOffset.from(">"); // ">" 表示未读消息 StreamReadOptions options = StreamReadOptions.empty().count(20).block(Duration.ofSeconds(2)); List<streammessage string>> messages = redisClient.read(groupName, consumerName, options, StreamOffset.create("order_events", offset));</streammessage>
确认与容错:手动 ACK + XPENDING 处理失败消息
消息被读出后进入 Pending Entries List(PEL),只有显式 ACK 才真正移除。宕机、超时、处理异常时,消息仍保留在 PEL 中,可被其他实例或重试逻辑捞起。
- 成功处理后,立即调用
XACK order_events pay_group <messageid></messageid> - 定期用
XPENDING order_events pay_group - + 10查看卡在 PEL 超过阈值(如 60 秒)的消息 - 对超时消息,可用
XCLAIM将其转移给当前消费者重新处理,或转入死信流(如order_events_dlq) - Spring Boot 中建议封装成 AOP 切面或 try-catch 后置确认,确保 ACK 不被遗漏
运维与清理:控制长度 + 监控 Pending
Stream 默认无限增长,需主动治理;Pending 过多说明消费瓶颈或 ACK 遗漏,是关键健康指标。
- 写入时加
MAXLEN ~ 1000000限制最大长度(如XADD order_events MAXLEN ~ 1000000 * field value),Redis 自动淘汰旧消息 - 用
XINFO GROUPS order_events查看各组 pending 数、consumer 数、group lag(落后消息数) - 对长期无人消费的消费者,用
XGROUP DELCONSUMER清理其 PEL,防止内存泄漏 - 生产环境建议暴露
pending-count指标到 Prometheus,设置告警阈值(如 > 1000)
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南










