Webflux结合Websocket场景下如何避免reactive redis消息操作重复订阅
问题根因
你的代码触发重复订阅有两个核心原因:
ReactiveRedisOperations.listenTo()方法每调用一次就会新建一个独立的Redis订阅,默认不会复用已有订阅- Websocket Handler中用
flatMap处理客户端请求,每次客户端发消息都会触发一次新的订阅创建,且旧订阅不会被主动取消,最终出现同一会话下多个订阅共存、重复收到消息的问题
修复方案
第一步:Service层实现全局订阅复用
调整MessagingService,将Redis订阅的Flux改造为热发布流,所有订阅者共享同一份上游Redis订阅,避免重复创建Redis连接:
public class MessagingService{ private final ReactiveRedisOperations<String, GamePubSub> reactiveRedisOperations; // 缓存全局共享的Topic订阅流,share()操作符实现多订阅者复用同一份上游订阅 // 当所有订阅者都断开时会自动取消Redis订阅,不占用多余资源 private final Flux<Object> sharedTopicFlux; public MessagingService(ReactiveRedisOperations<String, GamePubSub> reactiveRedisOperations) { this.reactiveRedisOperations = reactiveRedisOperations; this.sharedTopicFlux = reactiveRedisOperations.listenTo("TOPIC_NAME") .map(channelMessage -> channelMessage.getMessage()) // 可根据业务需求调整返回的字段 .share(); } public Flux<Object> playGame(UserInput userInput){ // 如果需要根据用户输入动态监听不同Topic,可维护一个ConcurrentHashMap<String, Flux<Object>> // 以Topic为key缓存对应订阅流,不存在时新建,已存在直接返回即可 return sharedTopicFlux; } }
第二步:Handler层避免同一会话重复订阅
将flatMap替换为switchMap,每次收到新的客户端消息时,自动取消之前的内部流订阅,确保同一会话最多只有一个活跃订阅:
public class SendingMessageHandler implements WebSocketHandler { private final Gson gson = new Gson(); private final MessagingService messagingService; public SendingMessageHandler(MessagingService messagingService) { this.messagingService = messagingService; } @Override public Mono<Void> handle(WebSocketSession session) { Flux<WebSocketMessage> stringFlux = session.receive() .map(WebSocketMessage::getPayloadAsText) // 此处补充你自己的inputData转UserInput逻辑 .switchMap(inputData -> messagingService.playGame(convertToUserInput(inputData)) .map(data -> session.textMessage(gson.toJson(data))) ); return session.send(stringFlux); } // 示例方法,根据自己的业务逻辑实现反序列化 private UserInput convertToUserInput(String input) { return gson.fromJson(input, UserInput.class); } }
扩展调整建议
- 如果需要给新连接推送最近N条历史消息,可以将
share()替换为replay(N).autoConnect(),N为需要缓存的历史消息条数 - 如果需要支持按用户过滤消息,拿到订阅流的消息后根据用户标识过滤即可,不要为不同用户创建重复的Topic订阅
内容的提问来源于stack exchange,提问作者mattts
相关产品推荐
相关产品推荐

