You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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
  • expiredKeySubscription Bean会在应用启动时自动完成事件订阅,返回的Disposable可用于在需要时手动取消订阅

内容的提问来源于stack exchange,提问作者JinYoung Park

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 09:50:52