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

Ktor WebSocket发送数据后如何同步获取客户端返回的String结果

Ktor WebSocket发送消息后同步获取客户端返回值实现方案

WebSocket本身是全双工异步协议,和HTTP的请求-响应模型不同,Ktor提供的DefaultWebSocketSession.send()方法仅负责出站消息发送,不会自动关联后续入站消息,因此返回值为Unit,需要自行实现请求-响应的关联逻辑,推荐用协程的CompletableDeferred实现非阻塞的同步等待。

实现思路

  • 给每一条待发送的消息生成唯一请求ID,随消息一起发给客户端,同时约定客户端回传消息时携带对应的请求ID
  • 为每个WebSocket会话维护待响应请求的映射表,key为请求ID,value为用来承载返回结果的CompletableDeferred对象
  • 启动独立协程持续监听当前会话的入站消息,收到客户端消息后解析出请求ID,找到对应的Deferred对象填充返回值
  • 发送消息后调用Deferred的await()方法,即可在当前协程上下文里挂起等待,直到拿到客户端返回的结果,整个过程不会阻塞线程。

可直接复用的代码示例

import io.ktor.websocket.*
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.channels.consumeEach
import kotlinx.coroutines.launch
import kotlinx.coroutines.withTimeout
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.Json
import java.util.UUID
import java.util.concurrent.TimeUnit

// WebSocket消息结构体,用于做请求响应关联
@Serializable
data class WsTransferMessage(
    val requestId: String,
    val payload: String
)

class KtorWsConnection(private val session: DefaultWebSocketSession) {
    // 待响应请求存储
    private val pendingResponseMap = mutableMapOf<String, CompletableDeferred<String>>()
    private val jsonParser = Json { ignoreUnknownKeys = true; encodeDefaults = true }
    // 响应超时时间,可根据业务调整
    private val defaultTimeoutMs = TimeUnit.SECONDS.toMillis(15)

    init {
        // 启动入站消息监听协程
        session.launch {
            session.incoming.consumeEach { frame ->
                if (frame !is Frame.Text) return@consumeEach
                val receivedContent = frame.readText()
                runCatching {
                    jsonParser.decodeFromString(WsTransferMessage.serializer(), receivedContent)
                }.onSuccess { msg ->
                    // 匹配到对应请求时填充返回结果
                    pendingResponseMap.remove(msg.requestId)?.complete(msg.payload)
                }
            }
            // 会话关闭时清空所有待处理请求,避免协程泄漏
            pendingResponseMap.values.forEach { it.cancel() }
            pendingResponseMap.clear()
        }
    }

    /**
     * 发送消息并等待客户端返回结果
     * @param data 待发送的业务数据
     * @return 客户端回传的字符串结果
     * @throws TimeoutCancellationException 超时未返回时抛出
     */
    suspend fun sendAndGetResponse(data: String): String {
        val requestId = UUID.randomUUID().toString()
        val deferred = CompletableDeferred<String>()
        pendingResponseMap[requestId] = deferred

        // 发送带requestId的消息
        val sendMsg = WsTransferMessage(requestId = requestId, payload = data)
        session.send(Frame.Text(jsonParser.encodeToString(WsTransferMessage.serializer(), sendMsg)))

        // 带超时等待返回
        return withTimeout(defaultTimeoutMs) {
            deferred.await()
        }
    }
}

简化场景适配

如果业务不存在并发发送消息的场景,同一时间只会发一条消息等返回,不需要支持多请求并发,可以去掉requestId和映射表,仅用一个全局的Deferred变量即可,逻辑会更简单。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 18:27:35