Spring Boot使用StreamMessageListenerContainer消费Redis重启后断连如何配置?
问题根因
默认配置下的StreamMessageListenerContainer没有开启连接异常后的自动订阅恢复机制,当Redis服务重启导致连接中断时,监听线程抛出异常后会直接终止消费任务,不会自动重建连接、重新注册监听。
解决步骤
1. 调整监听容器配置
在构建StreamMessageListenerContainerOptions时新增两个配置:
- 开启
resumeSubscriptionOnError(true),触发异常后自动恢复订阅 - 自定义错误处理器,记录异常日志便于排查问题
2. (可选)优化连接工厂配置
如果使用Lettuce作为Redis客户端,建议给连接工厂补充客户端重试、心跳检测配置,避免连接断开后无法正常重建。
修改后的代码示例
@Bean public Subscription listener(RedisStreamConsumer streamListener, RedisConnectionFactory redisConnectionFactory) throws InterruptedException { StreamMessageListenerContainer<String, MapRecord<String, Object, String>> listenerContainer = StreamMessageListenerContainer.create(redisTemplate().getConnectionFactory(), StreamMessageListenerContainer.StreamMessageListenerContainerOptions.builder() .hashKeySerializer(new StringRedisSerializer()).hashValueSerializer(new StringRedisSerializer()) .pollTimeout(Duration.ofMillis(100)) // 新增:异常后自动恢复订阅 .resumeSubscriptionOnError(true) // 新增:自定义错误处理器,记录异常 .errorHandler(ex -> { log.error("Redis Stream消费异常,将自动尝试恢复: ", ex); }) .build()); Subscription subscription = listenerContainer.receive(Consumer.from(groupName, consumerName), StreamOffset.create(consumerstreamName, ReadOffset.lastConsumed()), streamListener); subscription.await(Duration.ofSeconds(2)); listenerContainer.start(); return subscription; }
注意事项
- 如果你使用的Spring Data Redis版本低于2.6,没有
resumeSubscriptionOnError配置,需要手动实现异常监听逻辑:在监听器抛出连接异常时,调用listenerContainer.remove(subscription)删除失效订阅,再重新调用receive方法注册新的监听即可。 - 确保Redis服务开启了RDB或AOF持久化,避免重启后消费者组的消费偏移量丢失,导致重复消费或者从队列头重新消费的问题。
内容的提问来源于stack exchange,提问作者raghu babu parmala
相关产品推荐
相关产品推荐

