如何在Spring Data Reactive Redis中实现键过期事件监听器?
在Spring Data Reactive Redis中实现Key过期监听
同步版Spring Data Redis实现方式
Spring Data Redis可按如下方式实现Key过期监听:
@Component class ExpirationListener : MessageListener { override fun onMessage(message: Message, pattern: ByteArray) { println("########## onMessage pattern " + String(pattern) + " | " + message.toString()) } } @Configuration class RedisConfig(private val env: Environment) { @Bean fun redisConnectionFactory(): RedisConnectionFactory { val redisHost = env.getProperty("CONF_TRANSACTION_GATEWAY_REDIS_DB_HOST", "localhost") val redisPort = env.getProperty("CONF_TRANSACTION_GATEWAY_REDIS_DB_PORT", "6379") return LettuceConnectionFactory(redisHost, redisPort.toInt()) } @Bean fun redisTemplate(redisConnectionFactory: RedisConnectionFactory): RedisTemplate<String, String> { val stringSerializer = StringRedisSerializer() return RedisTemplate<String, String>() .apply { connectionFactory = redisConnectionFactory keySerializer = stringSerializer hashKeySerializer = stringSerializer valueSerializer = stringSerializer hashValueSerializer = stringSerializer } } @Bean fun redisMessageListenerContainer( redisConnectionFactory: RedisConnectionFactory, expirationListener: ExpirationListener ): RedisMessageListenerContainer { val redisMessageListenerContainer = RedisMessageListenerContainer() redisMessageListenerContainer.connectionFactory = redisConnectionFactory redisMessageListenerContainer.addMessageListener(expirationListener, PatternTopic("__keyevent@*__:expired")) return redisMessageListenerContainer } }
问题说明
但在Spring Data Reactive Redis中无法直接复用上述实现,核心问题是:
- 响应式场景下没有
RedisMessageListenerContainer类,更不存在addMessageListener()方法 - 同步的
MessageListener接口不适用响应式编程模型
响应式实现方案
在Spring Data Reactive Redis中,需通过ReactiveRedisConnection结合Reactor流实现Key过期事件监听,具体代码如下:
1. 定义响应式过期事件处理器
@Component class ReactiveExpirationHandler { fun handleExpiredKey(message: ByteBuffer) { val expiredKey = StandardCharsets.UTF_8.decode(message).toString() println("########## 响应式监听:Key已过期 $expiredKey") } }
2. 配置响应式Redis连接与事件订阅
@Configuration class ReactiveRedisConfig(private val env: Environment) { @Bean fun reactiveRedisConnectionFactory(): ReactiveRedisConnectionFactory { val redisHost = env.getProperty("CONF_TRANSACTION_GATEWAY_REDIS_DB_HOST", "localhost") val redisPort = env.getProperty("CONF_TRANSACTION_GATEWAY_REDIS_DB_PORT", "6379") return LettuceConnectionFactory(redisHost, redisPort.toInt()) } @Bean fun reactiveRedisTemplate(reactiveRedisConnectionFactory: ReactiveRedisConnectionFactory): ReactiveRedisTemplate<String, String> { val stringSerializer = StringRedisSerializer() val serializationContext = RedisSerializationContext .newSerializationContext<String, String>() .key(stringSerializer) .value(stringSerializer) .hashKey(stringSerializer) .hashValue(stringSerializer) .build() return ReactiveRedisTemplate(reactiveRedisConnectionFactory, serializationContext) } @Bean fun expiredKeySubscription( reactiveRedisConnectionFactory: ReactiveRedisConnectionFactory, expirationHandler: ReactiveExpirationHandler ): Disposable { // 订阅Key过期事件流 return reactiveRedisConnectionFactory.reactiveConnection .listenTo(PatternTopic("__keyevent@*__:expired")) .subscribe { message -> expirationHandler.handleExpiredKey(message.body) } } }
关键说明
- 使用
ReactiveRedisConnection.listenTo()方法订阅指定的PatternTopic,该方法返回Flux<ReactiveRedisMessage>,完全契合响应式流模型 - 通过
subscribe()方法处理每个过期事件,替代了同步版的MessageListener expiredKeySubscriptionBean会在应用启动时自动完成事件订阅,返回的Disposable可用于在需要时手动取消订阅
内容的提问来源于stack exchange,提问作者JinYoung Park
相关产品推荐
相关产品推荐

