Spring订阅Redis集群Key过期事件失败求助
环境信息
- Mac Ventura 13.6.3
- Temurin 17
- SpringBoot 3.2.1
- 依赖:
org.springframework.boot:spring-boot-starter-data-redis - 本地Redis集群(localhost:7001、localhost:7002、localhost:7003)
需求
实现通过Spring接收Redis集群的Key过期事件。
已执行操作
连接集群后执行以下Redis命令:
CONFIG SET notify-keyspace-events Ex SET hi 123 EXPIRE hi 3
Java代码
配置类
package kr.co.yogiyo.payo.infrastructure.temporal; import io.lettuce.core.ReadFrom; import kr.co.yogiyo.payo.infrastructure.temporal.service.RedisExpireEventService; import lombok.RequiredArgsConstructor; import org.springframework.boot.autoconfigure.data.redis.RedisConnectionDetails; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.redis.connection.RedisClusterConfiguration; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.connection.RedisNode; import org.springframework.data.redis.connection.lettuce.LettuceClientConfiguration; import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; import org.springframework.data.redis.listener.PatternTopic; import org.springframework.data.redis.listener.RedisMessageListenerContainer; @RequiredArgsConstructor @Configuration public class RedisConfig { private final RedisConnectionDetails redisConnectionDetails; @Bean public RedisConnectionFactory redisConnectionFactory() { RedisConnectionDetails.Cluster cluster = redisConnectionDetails.getCluster(); RedisClusterConfiguration clusterConfiguration = new RedisClusterConfiguration(); for (RedisConnectionDetails.Node node : cluster.getNodes()) { clusterConfiguration.addClusterNode(new RedisNode(node.host(), node.port())); } LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder().readFrom(ReadFrom.REPLICA_PREFERRED).build(); return new LettuceConnectionFactory(clusterConfiguration, clientConfig); } @Bean RedisMessageListenerContainer keyExpirationListenerContainer( RedisConnectionFactory connectionFactory, RedisExpireEventService redisExpireEventService) { RedisMessageListenerContainer listenerContainer = new RedisMessageListenerContainer(); listenerContainer.setConnectionFactory(connectionFactory); listenerContainer.addMessageListener( redisExpireEventService, new PatternTopic("__keyevent@*__:expired")); return listenerContainer; } }
服务类
package kr.co.yogiyo.payo.infrastructure.temporal.service; import lombok.RequiredArgsConstructor; import org.springframework.data.redis.connection.Message; import org.springframework.data.redis.connection.MessageListener; import org.springframework.stereotype.Service; @Service @RequiredArgsConstructor public class RedisExpireEventService implements MessageListener { @Override public void onMessage(Message message, byte[] pattern) { System.out.println("Message received: " + message.toString()); } }
application.yml
spring: main: banner-mode: off allow-bean-definition-overriding: true jackson: property-naming-strategy: SNAKE_CASE data: redis: host: ${PAYO_REDIS_MASTER_HOST:localhost} port: ${PAYO_REDIS_MASTER_PORT:7002} cluster: nodes: ${PAYO_REDIS_REPLICATION_NODES:localhost:7001,localhost:7002,localhost:7003}
问题
Spring项目可正常启动,但Key过期时无任何响应(调试模式下断点也未触发)。
可正常运行的Python代码
import redis def main(): # 连接Redis r = redis.Redis(host='localhost', port=7002, db=0) # 订阅__keyevent@0__:expired频道 pubsub = r.pubsub() pubsub.psubscribe('__keyevent@*__:expired') # 开始监听消息 for message in pubsub.listen(): print("有键过期! ", message) if __name__ == "__main__": main()
运行输出
有键过期! {'type': 'psubscribe', 'pattern': None, 'channel': b'__keyevent@*__:expired', 'data': 1} 有键过期! {'type': 'pmessage', 'pattern': b'__keyevent@*__:expired', 'channel': b'__keyevent@0__:expired', 'data': b'hi'}
提问
我的Java代码中缺少了什么配置或逻辑?
问题排查与解决方案
全节点配置notify-keyspace-events
Redis集群中每个节点独立运行,你只在单个节点执行了CONFIG SET notify-keyspace-events Ex,但Key可能存储在任意节点,过期事件只会在Key所在节点触发。需要给集群所有节点执行该命令;若要永久生效,需修改每个节点的redis.conf文件,添加notify-keyspace-events Ex后重启节点。配置Lettuce集群拓扑刷新
集群模式下,Lettuce默认仅在初始连接节点订阅频道,无法接收其他节点的过期事件。需添加拓扑刷新配置,让客户端自动发现集群节点并全节点订阅:
修改redisConnectionFactory方法中的Lettuce客户端配置:import io.lettuce.core.ClientOptions; import io.lettuce.core.cluster.ClusterTopologyRefreshOptions; // ... LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder() .readFrom(ReadFrom.REPLICA_PREFERRED) .clientOptions(ClientOptions.builder() .clusterTopologyRefreshOptions(ClusterTopologyRefreshOptions.builder() .enableAllAdaptiveRefreshTriggers() .build()) .build()) .build();指定Redis数据库编号
Python代码连接的是db=0,检查Spring是否使用相同数据库,若默认不是db=0,需在application.yml中添加:spring: data: redis: database: 0订阅专用连接工厂(可选)
过期事件仅在主节点触发,REPLICA_PREFERRED策略可能导致订阅连接到副本节点无法接收事件。可给订阅容器单独配置主节点优先的连接工厂:@Bean RedisConnectionFactory redisSubscriptionConnectionFactory() { RedisConnectionDetails.Cluster cluster = redisConnectionDetails.getCluster(); RedisClusterConfiguration clusterConfiguration = new RedisClusterConfiguration(); for (RedisConnectionDetails.Node node : cluster.getNodes()) { clusterConfiguration.addClusterNode(new RedisNode(node.host(), node.port())); } LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder() .readFrom(ReadFrom.MASTER) .clientOptions(ClientOptions.builder() .clusterTopologyRefreshOptions(ClusterTopologyRefreshOptions.builder() .enableAllAdaptiveRefreshTriggers() .build()) .build()) .build(); return new LettuceConnectionFactory(clusterConfiguration, clientConfig); } @Bean RedisMessageListenerContainer keyExpirationListenerContainer( RedisExpireEventService redisExpireEventService) { RedisMessageListenerContainer listenerContainer = new RedisMessageListenerContainer(); listenerContainer.setConnectionFactory(redisSubscriptionConnectionFactory()); listenerContainer.addMessageListener( redisExpireEventService, new PatternTopic("__keyevent@*__:expired")); return listenerContainer; }
内容的提问来源于stack exchange,提问作者Jerry
相关产品推荐
相关产品推荐

