Spring WebFlux WebSocket跨服务器请求代理实现问题求助
问题解决思路与代码修正
核心问题分析
你的代码无法编译且不符合WebFlux响应式编程模型:
process方法声明返回同步的Response,但WebFlux的WebSocket操作是异步非阻塞的,不能直接同步返回结果- 调用
subscribe()会返回Disposable(用于取消订阅),而非业务响应数据,这是响应式编程的基础特性
代码修正方案
服务器A(WebSocket客户端)修改
将方法返回类型改为Mono<Response>,并调整响应式流的组合逻辑,确保先发送请求再接收响应:
private val webSocketClient: WebSocketClient = ReactorNettyWebSocketClient() fun process(request: Request): Mono<Response> { return webSocketClient.execute(URI.create("ws://localhost:8080/websocket")) { session -> // 构造并发送请求消息 val sendRequest = session.send( Mono.just(request) .map { objectMapper.writeValueAsString(it) } .map { session.textMessage(it) } ) // 接收响应消息,取第一个结果(匹配请求-响应的单次交互场景) val receiveResponse = session.receive() .map { objectMapper.readValue<Response>(it.payloadAsText) } .next() // 先完成发送操作,再返回响应结果 sendRequest.then(receiveResponse) } }
关键调整说明
- 返回
Mono<Response>:符合WebFlux异步编程模型,由上层调用方(如HTTP控制器)负责订阅触发执行,避免线程阻塞 sendRequest.then(receiveResponse):确保请求发送完成后再等待响应,保证流的执行顺序next():针对请求-响应的单次交互场景,只取WebSocket返回的第一条消息作为业务响应
服务器B(WebSocket服务端)补充
你的现有代码逻辑基本正确,只需确保processMessage方法返回Response类型,保证流的类型匹配:
internal class WebSocketHandler( private val objectMapper: ObjectMapper, ) : WebSocketHandler { override fun handle(session: WebSocketSession): Mono<Void> { val outputMessages = session .receive() .map { objectMapper.readValue<Request>(it.payloadAsText) } .map { processMessage(it) } // 确保此方法返回Response类型 .map { objectMapper.writeValueAsString(it) } .map { session.textMessage(it) } return session.send(outputMessages) } private fun processMessage(request: Request): Response { // 实现你的业务处理逻辑,返回Response } }
概念层面提示
- 不要在WebFlux代码中尝试同步阻塞获取结果,响应式编程的核心是通过
Mono/Flux传递异步操作流 - WebSocket是全双工协议,但你的场景是请求-响应模式,因此用
next()提取单条响应即可满足需求
内容的提问来源于stack exchange,提问作者stosik
相关产品推荐
相关产品推荐

