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

Spring WebFlux WebSocket跨服务器请求代理实现问题求助

问题解决思路与代码修正

核心问题分析

你的代码无法编译且不符合WebFlux响应式编程模型:

  1. process方法声明返回同步的Response,但WebFlux的WebSocket操作是异步非阻塞的,不能直接同步返回结果
  2. 调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 11:17:53