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

