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

Webflux结合Websocket场景下如何避免reactive redis消息操作重复订阅

问题根因

你的代码触发重复订阅有两个核心原因:

  1. ReactiveRedisOperations.listenTo()方法每调用一次就会新建一个独立的Redis订阅,默认不会复用已有订阅
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 22:15:08