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

Spring Boot中如何监听Redis缓存更新事件并返回Flux<Data>

实现Redis缓存更新时推送WebFlux的Flux

要实现监听指定Redis键的更新并返回Flux,核心是利用Redis的**键空间通知(Keyspace Notifications)**结合WebFlux的响应式能力,步骤如下:

1. 开启Redis键空间通知

Redis默认关闭键空间通知,需要先开启它才能监听键的更新事件:

  • 持久化配置:修改Redis配置文件redis.conf,添加/修改:
    notify-keyspace-events KEA
    
    其中K表示监听键空间事件,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 04:29:53