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

