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

如何让Spring WebFlux WebSocket与Kotlin协程兼容?

Spring WebFlux WebSocket 与 Kotlin 协程适配方案

原代码的问题在于doOnNext是Reactor的同步操作符,无法直接调用Kotlin挂起函数——挂起函数的协程调度逻辑不会被正确触发,导致代码无法按预期运行。以下是无需改用Reactor实现、纯协程风格的适配方案:

1. 确保依赖到位

首先确认引入kotlinx-coroutines-reactor库,它是协程与Reactor适配的核心:

// Gradle 依赖示例
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-reactor:1.7.3")
implementation("org.springframework.boot:spring-boot-starter-webflux")

2. 修改核心代码

用mono { }将挂起函数包装为Reactor流,替换原有的doOnNext,同时用协程风格的awaitSingle()替代阻塞式的block():

import kotlinx.coroutines.reactor.mono
import kotlinx.coroutines.runBlocking
import org.springframework.web.reactive.socket.WebSocketSession
import org.springframework.web.reactive.socket.client.ReactorNettyWebSocketClient
import java.net.URI

fun main() = runBlocking {
    val client = ReactorNettyWebSocketClient()
    val uri = URI.create("wss://websocket")

    client.execute(uri) { session ->
        val handler = WebSocketImplementation(session)
        session.receive()
            .flatMap { message ->
                // 用mono包装挂起函数,让协程逻辑融入Reactor流
                mono {
                    handler.handleTextMessage(message)
                    // 必须消费消息payload,避免资源泄漏
                    message.payload.receive().awaitFirstOrNull()
                }
            }
            .then()
    }.awaitSingle()
}

class WebSocketImplementation(private val session: WebSocketSession) {
    suspend fun handleTextMessage(message: org.springframework.web.reactive.socket.WebSocketMessage) {
        // 你的协程逻辑示例:处理消息并发送响应
        val processedContent = "Received: ${message.payloadAsText}"
        session.send(session.textMessage(processedContent)).awaitFirstOrNull()
    }
}

关键说明

  • mono { }的作用:将挂起函数的协程逻辑包装为Reactor的Mono类型,确保协程的挂起、恢复被正确调度,Reactor流会等待挂起函数执行完成后再推进。
  • 替换block()为awaitSingle():在runBlocking协程作用域内使用非阻塞的协程等待方法,既保持协程风格,又避免线程阻塞。
  • 消费消息payload:WebSocket消息的payload必须被消费才能释放底层资源,添加message.payload.receive().awaitFirstOrNull()确保资源正常回收。
  • 异常处理:如果需要捕获挂起函数的异常,可以在mono { }内部添加try-catch,或者在Reactor流上追加onErrorResume等操作符处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 23:20:05