必须先手动创建消费者组,否则XREADGROUP直接报错;正确做法是用XGROUP CREATE mystream mygroup $ MKSTREAM显式初始化,其中$表示从新消息开始消费,MKSTREAM自动创建Stream。

必须先手动创建消费者组,否则XREADGROUP直接报错
Spring Boot 的 StringRedisTemplate 不会自动帮你建消费者组,哪怕 Stream 已存在。调用 XREADGROUP 时如果组不存在,Redis 会返回 NOGROUP No such consumer group 错误,而不是静默创建。
正确做法是在应用启动时或首次消费前,显式执行初始化命令:
XGROUP CREATE mystream mygroup $ MKSTREAM
其中 $ 表示从最新消息开始消费(即只处理后续新进消息),MKSTREAM 表示若 Stream 不存在则自动创建。这个操作只需执行一次,建议封装在 @PostConstruct 方法或 Spring Boot 的 ApplicationRunner 中。
消费者必须显式调用XACK,否则消息永远卡在PENDING列表
Stream 的可靠性依赖于消费确认机制。只要没调用 XACK,那条消息就会一直留在该消费者组的 PENDING 列表里,下次 XREADGROUP 还会再次读到它——这既是“至少一次”语义的保障,也是容易被忽略的资源泄漏点。
常见错误写法是只读不确认,或在 try-catch 外层漏掉 XACK:
- 务必在业务逻辑成功执行后,立即调用
stringRedisTemplate.opsForStream().acknowledge("mystream", "mygroup", record.getId()) - 如果处理失败需重试,不要
XACK;可稍后通过XCLAIM抢回超时消息,但生产环境更推荐直接抛异常触发重入队逻辑 - 注意:
record.getId()是字符串类型,如"1725562800000-0",不能传 null 或空字符串
StringRedisTemplate中XADD必须用"*"作为ID,且消息体建议用JSON字符串
StringRedisTemplate 对 XADD 的封装比较严格:若传入 null 作为消息 ID,会直接抛 IllegalArgumentException;而传 "*" 才能触发 Redis 自动生成唯一 ID(格式为 毫秒时间戳-序列号)。
消息体也建议统一用 JSON 字符串,而非 Java 对象序列化:
- 避免不同服务间 class 版本不一致导致反序列化失败
- 便于用
redis-cli直接调试和人工补数据 - 示例写法:
StreamRecords.newRecord().in("mystream").withFields(Map.of("event", "order_created", "order_id", "1001"))
单次XREADGROUP的COUNT不宜过大,尤其在高吞吐场景
XREADGROUP 的 COUNT 参数控制每次拉取的消息数量。设为 100 甚至 1000 看似能提升吞吐,但实际容易引发问题:
- 单次处理耗时变长,可能触发超时或 OOM(尤其消息体较大时)
- 若中途崩溃,未
XACK的消息批量堆积在 PENDING 列表,影响后续消费进度 - 推荐值为
1~10,配合短周期轮询(如BLOCK 5000),兼顾响应性与吞吐 - 若真需批量处理,应在业务层做聚合,而非依赖 Redis 一次吐出大量记录
最容易被忽略的是消费者组生命周期管理:没有类似 Kafka 的自动 rebalance,组内消费者宕机后,其 PENDING 消息默认不会自动转移——得靠 XCLAIM 或设置 TIMEOUT 配合定时任务来兜底。这点在灰度发布或滚动更新时特别关键。


















