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
相关产品推荐
相关产品推荐

