spring boot 3 中用 redis stream 做消息队列必须显式创建消费者组、手动调用 xack、主动处理 pending 列表,否则消息会堆积、重复或丢失;这是设计契约而非配置问题。

直接说结论:Spring Boot 3 中用 Redis Stream 做消息队列,必须显式创建消费者组、手动调用 XACK、主动处理 PENDING 列表,否则消息会堆积不消费、重复投递或永久丢失——这不是配置问题,是设计契约。
必须先执行 XGROUP CREATE,否则 XREADGROUP 报 NOGROUP
Spring Boot 不会自动帮你建消费者组,XREADGROUP 第一次调用就失败是高频报错。错误信息是:NOGROUP No such consumer group。
- 不能依赖应用启动时“顺手”创建:Stream 可能还没写入首条消息,
XGROUP CREATE ... MKSTREAM就会失败;必须确保 Stream 已存在(比如先发一条空消息,或用DEL mystream+XADD mystream * dummy ""初始化) - 推荐在初始化 Bean 阶段用
StringRedisTemplate.opsForStream().createGroup()显式调用,传参注意:groupName是字符串,readOffset推荐用$(从最新开始),不是0或空字符串 - 如果多个服务共用同一组名,要确认是否真需要共享消费进度;否则应按服务名区分组,避免互相干扰
消费者必须显式调用 XACK,否则消息永远在 PENDING 里
XREADGROUP 拉到消息后,Redis 就把它放进该消费者专属的 PENDING 列表,直到你调用 XACK。不调?下次还给你,且越积越多。
Redis 缓存和数据结构管理技能。通过自然语言操作 Redis,支持 String、Hash、List、Set、ZSet、Stream 等数据结构操作。当用户提到 Redis、缓存、消息队列、会话存储时使用此技能。
- Spring Data Redis 的
StreamOperations.readGroup()返回的是Map<recordid map object>></recordid>,但不会自动 ACK;你得自己遍历结果,对每条成功处理的消息调用acknowledge() - 异常路径必须兜底:try-catch 里不能只 log,要在 finally 或 catch 块中判断是否已处理,未处理则跳过 ACK;否则异常中断导致没 ACK,消息卡死
- 别信“自动 ACK”配置:Spring Boot 3.x 的
StreamMessageListenerContainer默认不 ACK,也没有开关;它只是帮你拉数据,ACK 是你的责任
消息体用 JSON 字符串,别序列化 Java 对象
用 Object 直接塞进 StreamRecords.newRecord().withObject() 看似方便,但反序列化时极易因 class 版本不一致、字段缺失、JDK 升级崩溃。
- 生产者侧统一走
objectMapper.writeValueAsString(event),再传给withFieldObject("data", jsonStr) - 消费者侧拿到
Map<string object></string>后,取"data"字段再用objectMapper.readValue(jsonStr, Event.class) - 字段命名建议小驼峰(如
orderId),避免不同语言客户端(如 Go/Python 消费者)解析歧义 - 如果消息体含二进制(如图片 base64),务必额外加
content-type字段标注,否则下游无法安全 decode
PENDING 消息要定期扫描,否则故障后无法自愈
消费者宕机、OOM 或网络分区后,PENDING 里的消息不会自动转移。Redis 不会主动“超时重分配”,得你自己查、自己 XCLAIM。
- 用定时任务(如
@Scheduled(fixedDelay = 30000))调用pendingMessages()扫描每个消费者组的 pending 列表 - 筛选条件必须带
minIdleTime(如 60_000L),只处理闲置超 1 分钟的消息,避免刚取出来就误判为失败 -
XCLAIM时指定新 consumer 名(如"recovery-worker"),并设idle为 0,防止二次抢占 - 别漏日志:每次
XCLAIM成功后记一条 warn 日志,包含原 consumer、ID、idle 时间,这是线上排查消费卡顿的第一线索
最易被忽略的一点:Redis Stream 的可靠性不来自某个 API 开关,而来自你对 PENDING 列表的持续治理——它不像 Kafka 有 coordinator 自动 rebalance,也不像 RabbitMQ 有 connection timeout 触发 redeliver。你写的每一行 ACK 和 CLAIM 代码,都在定义这条消息的生死边界。










