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

Kotlin中ReactorNettyWebSocketClient无法使用subscribe/subscribeWith问题求助

解决Kotlin中ReactorNettyWebSocketClient无法使用subscribe/subscribeWith的问题

我明白你遇到的困扰了——ReactorNettyWebSocketClient的execute方法返回的是Mono<Void>,它代表整个WebSocket会话的生命周期,而非可直接订阅的消息流,所以没法直接在它上面调用subscribe/subscribeWith来处理消息。不过我们可以通过正确处理会话内的流逻辑,再订阅会话生命周期Mono来启动整个流程。

下面是修正后的完整代码,既能成功连接服务器、发送消息,还能正确处理返回的消息流,同时解决订阅的问题:

import reactor.netty.http.client.ReactorNettyWebSocketClient
import java.net.URI
import reactor.core.publisher.Flux
import org.springframework.web.reactive.socket.WebSocketSession
import org.springframework.web.reactive.socket.WebSocketMessage

fun main(args: Array<String>) {
    val uri = URI("ws://localhost:8080/myservice")
    val client = ReactorNettyWebSocketClient()
    // 替换为你的实际消息获取逻辑
    val input = Flux.just(readMsg())

    // 核心:订阅execute返回的Mono来启动整个WebSocket会话
    client.execute(uri) { session: WebSocketSession ->
        // 构建发送消息的流
        val sendFlow = session.send(input.map(session::textMessage))
        // 构建接收并处理消息的流
        val receiveFlow = session.receive()
            .map(WebSocketMessage::getPayloadAsText)
            .doOnNext { receivedMsg -> println("收到服务器消息:$receivedMsg") }
            .then() // 标记接收流处理完成

        // 合并发送与接收流,等待两者都完成后结束会话
        sendFlow.then(receiveFlow)
    }.doOnError { error -> println("连接出错:${error.message}") }
        .block() // 测试场景用block阻塞主线程,避免程序直接退出;生产环境可改用subscribe()异步启动
}

// 你的消息读取方法示例
fun readMsg(): String {
    return "客户端发送的测试消息"
}

关键要点说明:

  • execute的返回值逻辑:它返回的Mono<Void>是整个WebSocket会话的生命周期信号,只有当会话正常关闭或异常终止时,这个Mono才会完成。我们必须订阅它来触发整个连接流程。
  • 会话内的流组合:在execute的lambda中,需要将发送流和接收流通过then组合,确保会话会等待收发逻辑都完成后再关闭。
  • 避免单独订阅接收流:不需要在receive()返回的Flux上单独调用subscribe,doOnNext已经可以处理每条消息,整个流的订阅由execute返回的Mono驱动,WebFlux会自动绑定流的生命周期与会话。

内容的提问来源于stack exchange,提问作者Garth Gilmour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:30:30