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

Spring中如何向同一WebSocket会话连续推送多条消息并保持会话开启

在Spring中向同一个WebSocket Session连续推送消息并保持会话开启

要实现向同一个WebSocket Session连续推送多条消息且保持会话活跃,核心是确保发送流和接收流持续绑定到Session上,避免Session因发送完成而关闭。以下是具体说明和修正方案:

核心思路

Spring WebFlux的WebSocketClient通过execute方法建立会话,会话的生命周期由传入的Function<WebSocketSession, Mono<Void>>控制:只有当这个Mono完成时,Session才会关闭。因此,我们需要将发送消息的流和接收消息的流组合成一个持续的流,让Session一直处于活跃状态。

现有代码问题分析

  • JSON格式错误:input和input1中的JSON字符串多了一个左大括号,比如"{{\"action\":\"subscribe\""应改为"{\"action\":\"subscribe\""
  • 未利用所有发送流:input1和input2没有被加入到发送流程中
  • 会话过早关闭风险:当前代码中send(input...)完成后,后续的thenMany(receive())若没有持续消息,可能导致Mono完成、Session关闭
  • 不必要的延迟:delayUntil(i -> Mono.delay(Duration.ofSeconds(5)))会延迟每条接收消息的处理,可能影响会话活跃性
  • 返回值不合理:方法声明返回Flux<WebSocketSession>但最终返回null,不符合逻辑

修正后的实现代码

public Flux<String> openWebSocket() {
    // 修正JSON格式,合并所有需要发送的消息流
    Flux<String> messagesToSend = Flux.just(
            "{\"action\":\"auth\",\"params\":\"SomeKEY AUTH\"}\n",
            "{\"action\":\"subscribe\", \"params\":\"subcriptionmessage.*\"}\n",
            "{\"action\":\"subscribe\", \"params\":\"AM.*\"}\n",
            "{\"action\":\"auth\",\"params\":\"SomeKEY AUTH\"}\n"
    );

    ReplayProcessor<String> output = ReplayProcessor.create(100);

    try {
        webSocketClient.execute(new URI("wss://***"), session -> {
            log.debug("WebSocket会话已建立,开始发送消息");

            // 发送消息流:发送完成后保持流活跃,避免Session关闭
            Mono<Void> sendMono = session.send(
                    messagesToSend
                            .doOnNext(msg -> log.debug("发送出站消息: {}", msg))
                            .map(session::textMessage)
            ).then(Mono.never());

            // 接收消息流:持续处理服务器返回的消息
            Flux<String> receiveFlux = session.receive()
                    .map(WebSocketMessage::getPayloadAsText)
                    .doOnNext(msg -> log.debug("接收入站消息: {}", msg))
                    .subscribeWith(output);

            // 并行绑定发送和接收流,维持会话活跃
            return Mono.zip(sendMono, receiveFlux.then())
                    .then();
        }).subscribe();
    } catch (URISyntaxException e) {
        throw new RuntimeException("WebSocket URI格式错误", e);
    }

    return output;
}

关键细节说明

  • 合并发送流:将所有需要发送的消息合并到一个Flux中,实现连续推送
  • 保持发送流活跃:使用Mono.never()让发送流在完成后不终止,避免Session因发送完成而关闭
  • 并行处理流:通过Mono.zip让发送流和接收流并行执行,只要任一流处于活跃状态,Session就会保持开启
  • 正确返回结果:返回ReplayProcessor作为接收消息的输出流,调用方可订阅该流获取WebSocket的响应消息

内容的提问来源于stack exchange,提问作者Nirav Kumar Patel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 02:22:03