Spring Boot中如何监听Redis缓存更新事件并返回Flux<Data>
实现Redis缓存更新时推送WebFlux的Flux
要实现监听指定Redis键的更新并返回Flux,核心是利用Redis的**键空间通知(Keyspace Notifications)**结合WebFlux的响应式能力,步骤如下:
1. 开启Redis键空间通知
Redis默认关闭键空间通知,需要先开启它才能监听键的更新事件:
- 持久化配置:修改Redis配置文件
redis.conf,添加/修改:
其中notify-keyspace-events KEAK表示监听键空间事件,E表示监听键事件,A表示所有类型的事件。 - 临时生效:通过Redis CLI执行命令(重启Redis后失效):
CONFIG SET notify-keyspace-events KEA
2. 配置响应式Redis监听容器
在Spring Boot中配置ReactiveRedisMessageListenerContainer,用于响应式接收Redis通知:
@Configuration public class RedisConfig { @Bean public ReactiveRedisMessageListenerContainer reactiveRedisMessageListenerContainer(ReactiveRedisConnectionFactory connectionFactory) { return new ReactiveRedisMessageListenerContainer(connectionFactory); } // 配置Data对象的序列化,确保能正确读写Redis @Bean public ReactiveRedisOperations<String, Data> reactiveRedisOperations(ReactiveRedisConnectionFactory factory) { Jackson2JsonRedisSerializer<Data> dataSerializer = new Jackson2JsonRedisSerializer<>(Data.class); RedisSerializationContext<String, Data> serializationContext = RedisSerializationContext .<String, Data>newSerializationContext(new StringRedisSerializer()) .value(dataSerializer) .build(); return new ReactiveRedisTemplate<>(factory, serializationContext); } }
3. 实现缓存更新监听服务
创建服务类,监听指定键的SET事件(缓存更新/新增时通常触发该操作),返回更新后的Data对象Flux:
@Service public class CacheUpdateNotificationService { private final ReactiveRedisOperations<String, Data> reactiveRedisOps; private final ReactiveRedisMessageListenerContainer listenerContainer; public CacheUpdateNotificationService(ReactiveRedisOperations<String, Data> reactiveRedisOps, ReactiveRedisMessageListenerContainer listenerContainer) { this.reactiveRedisOps = reactiveRedisOps; this.listenerContainer = listenerContainer; } public Flux<Data> listenForKeyUpdates(String cacheKey) { // 构造键空间通知频道,格式为 __keyspace@<db编号>__:<键名>,默认db是0 String notificationChannel = "__keyspace@0__:" + cacheKey; return listenerContainer.receive(ChannelTopic.of(notificationChannel)) // 只过滤SET事件 .filter(message -> "SET".equals(message.getMessage())) // 获取更新后的缓存值 .flatMap(message -> reactiveRedisOps.opsForValue().get(cacheKey)) // 过滤空值,避免推送无效数据 .filter(Objects::nonNull); } }
4. 暴露订阅接口
在Controller中提供SSE(Server-Sent Events)接口,供客户端订阅缓存更新:
@RestController @RequestMapping("/api/cache") public class CacheUpdateController { private final CacheUpdateNotificationService notificationService; public CacheUpdateController(CacheUpdateNotificationService notificationService) { this.notificationService = notificationService; } @GetMapping(value = "/updates/{key}", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<Data> subscribeToCacheUpdates(@PathVariable String key) { return notificationService.listenForKeyUpdates(key); } }
注意事项
- 确保Spring Cache的Redis实现(如
RedisCacheManager)或自定义Redis操作在更新缓存时执行SET类命令,否则无法触发监听事件。 - 若Redis使用非默认数据库,需修改通知频道中的db编号。
内容的提问来源于stack exchange,提问作者Farfetch'd
相关产品推荐
相关产品推荐

