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

Spring Webflux ReactorNettyWebSocketClient接收数据时无法发送ping消息问题

问题根因
  • 核心原因是WebSocketSession.send()属于冷发布器,只有被订阅时才会真正执行发送逻辑。你在doOnNext回调中仅构造了发送动作实例,没有将其接入响应式流的订阅链路,因此Ping消息从未实际发送,你看到的日志仅为手动打印的副作用,和真实发送行为无关。
  • ReactorNettyWebSocketClient的execute方法要求会话的所有入站、出站逻辑必须整合进返回的Publisher生命周期中,单独调用send方法无法共享会话的出站通道上下文,自然无法完成消息发送。
解决方案

方案1:合并入站处理流与出站Ping流

将消息接收逻辑、定时Ping发送逻辑全部整合到会话的处理流中,由框架统一订阅执行:

// 有效token可从POST https://api.kucoin.com/api/v1/bullet-public接口的data.token字段获取
String token = "2neAiuYvAU61ZDXANAGAsiL4-iAExhsBXZxftpOeh_55i3Ysy2q2LEsEWU64mdzUOPusi34M_wGoSf7iNyEWJ_FJZkT3sFi"
            + "L0St-iwQclmJ_Sarx5w80udiYB9J6i9GjsxUuhPw3BlrzazF6ghq4LwBeRfXaJvyCxEfL2zsVCuw=.wjcFO2RaGwCtkQxvSvCqDA==";
URI wsUri = URI.create("wss://ws-api.kucoin.com/endpoint");
URI fullUri = UriComponentsBuilder.fromUri(wsUri).queryParam("token", token).build().toUri();

new ReactorNettyWebSocketClient().execute(fullUri, session -> {
    // 发送订阅请求
    Mono<Void> subscribeSender = session.send(Mono.just(
        session.textMessage("{\"id\":1,\"type\":\"subscribe\",\"topic\":\"/spotMarket/level2Depth5:BTC-USDT\"}")
    ));

    // 入站消息处理流,持续接收服务端推送
    Flux<Void> receiver = session.receive()
        .doOnNext(msg -> log.info("Received message: {}", msg.getPayloadAsText()))
        .then()
        .repeat(); // 保持接收流不终止

    // 定时Ping发送流,每10秒生成Ping消息交给会话发送
    Flux<Void> pingSender = Flux.interval(Duration.ofSeconds(10))
        .map(seq -> session.textMessage("{\"id\":\"1\",\"type\":\"ping\"}"))
        .doOnNext(msg -> log.info("Sent ping message"))
        .as(session::send);

    // 整合所有流,会话断开前保持运行
    return subscribeSender.thenMany(Flux.merge(receiver, pingSender)).then();
}).subscribe();

方案2:用Sinks实现异步消息发送

如果后续需要动态发送其他消息(比如新增订阅、取消订阅),可以用Sinks构造出站消息队列,所有待发送消息都推到队列中统一处理:

new ReactorNettyWebSocketClient().execute(fullUri, session -> {
    // 构造多播出站队列
    Sinks.Many<WebSocketMessage> outboundSink = Sinks.many().unicast().onBackpressureBuffer();
    
    // 订阅队列,统一发送所有出站消息
    Mono<Void> outboundSender = session.send(outboundSink.asFlux()).then();

    // 发送订阅消息
    outboundSink.tryEmitNext(session.textMessage("{\"id\":1,\"type\":\"subscribe\",\"topic\":\"/spotMarket/level2Depth5:BTC-USDT\"}"));

    // 定时推Ping消息到队列
    Disposable pingTask = Flux.interval(Duration.ofSeconds(10))
        .doOnNext(seq -> {
            outboundSink.tryEmitNext(session.textMessage("{\"id\":\"1\",\"type\":\"ping\"}"));
            log.info("Sent ping message");
        }).subscribe();

    // 入站处理流
    Flux<Void> receiver = session.receive()
        .doOnNext(msg -> log.info("Received message: {}", msg.getPayloadAsText()))
        .then()
        .repeat();

    // 会话结束后清理定时任务
    return Flux.merge(receiver, outboundSender)
        .doFinally(signalType -> pingTask.dispose())
        .then();
}).subscribe();
注意事项
  • KuCoin的WebSocket接口要求Ping消息的id字段为字符串类型,不要传数值类型,避免服务端识别失败导致连接断开。
  • 建议增加Pong消息处理逻辑,收到服务端返回的Pong后重置Ping定时器,降低网络波动带来的连接误断概率。

内容的提问来源于stack exchange,提问作者Dmytro D

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 14:36:05