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

基于Ktor WebSocket的MVVM架构聊天客户端实现疑问与最佳实践咨询

基于Ktor WebSocket的MVVM架构聊天客户端实现疑问与最佳实践咨询

看起来你已经在MVVM架构下用Ktor WebSocket搭建了聊天客户端的核心框架,这里针对你遇到的几个核心问题,结合Android/Kotlin的最佳实践来梳理下解决方案:


1. 在authenticate()中等待认证响应的优雅实现

你当前的authenticate()方法只发送了认证消息,但没有等待服务器的确认回复。我们可以结合StateFlow和挂起函数的特性,让authenticate()挂起直到收到认证结果:

步骤1:在Service中添加认证结果的StateFlow

class ChatSocketServiceImpl(private val tokenDataStore: TokenDataStore): ChatSocketService {
    // 新增:用于暂存认证结果的StateFlow
    private val _authResult = MutableStateFlow<Resource<Unit>?>(null)
    private val authResult: StateFlow<Resource<Unit>?> = _authResult

    // 原有状态变量改为StateFlow(后面会说明原因)
    private val _isConnected = MutableStateFlow(false)
    val isConnected: StateFlow<Boolean> = _isConnected
    private val _isAuthenticated = MutableStateFlow(false)
    val isAuthenticated: StateFlow<Boolean> = _isAuthenticated

    private val client = HttpClient(CIO) { install(WebSockets) }
    private var socket: DefaultClientWebSocketSession? = null

    // ... 其他原有代码
}

步骤2:在消息监听中处理认证结果

修改observeMessages()的map逻辑,当收到服务器的认证成功/失败消息时,同步更新状态和认证结果流:

override fun observeMessages(): Flow<ReceivedSocketMessage> {
    return try {
        socket?.incoming?.receiveAsFlow()
            ?.filter { it is Frame.Text }
            ?.map {
                val json = (it as? Frame.Text)?.readText() ?: ""
                val receivedMessageDto = Json.decodeFromString<ReceivedSocketMessage>(json)

                // 处理认证相关消息
                when (receivedMessageDto) {
                    is ReceivedSocketMessage.AuthenticationSuccess -> {
                        _isAuthenticated.value = true
                        _authResult.value = Resource.Success(Unit)
                    }
                    is ReceivedSocketMessage.AuthenticationFail -> {
                        _isAuthenticated.value = false
                        _authResult.value = Resource.Error("Authentication failed: ${receivedMessageDto.reason}")
                    }
                    else -> {} // 其他消息正常向下传递
                }

                Log.i("logs", "Received message from Chat Websocket: $receivedMessageDto")
                receivedMessageDto
            } ?: flow { /* 空流 */ }
    } catch(e: Exception) {
        e.printStackTrace()
        flow { }
    }
}

步骤3:改造authenticate()方法挂起等待结果

override suspend fun authenticate(): Resource<Unit> {
    return try {
        if (socket?.isActive != true) {
            return Resource.Error("No active socket connection")
        }

        _authResult.value = null // 重置之前的认证结果
        val authMessage = SendSocketMessage.Authenticate(
            token = tokenDataStore.getToken().toString()
        ).toString()
        socket?.outgoing?.send(Frame.Text(authMessage))

        // 挂起等待认证结果,添加超时避免无限等待
        withTimeout(5000) {
            authResult.first { it != null } ?: Resource.Error("No response from server")
        }
    } catch (e: TimeoutCancellationException) {
        _isAuthenticated.value = false
        Resource.Error("Authentication timed out")
    } catch (e: Exception) {
        _isAuthenticated.value = false
        e.printStackTrace()
        Resource.Error(e.localizedMessage ?: "Unknown error during authentication")
    }
}

2. 给initSession()和authenticate()添加重试机制

可以用Kotlin协程的retry扩展(需要引入依赖),或者自己实现指数退避的重试逻辑:

方式1:用官方retry扩展

首先在build.gradle添加依赖:

implementation "org.jetbrains.kotlinx:kotlinx-coroutines-retry:1.0.0"

然后改造initSession():

override suspend fun initSession(retryCount: Int = 3, initialDelayMs: Long = 1000): Resource<Unit> {
    return try {
        // 先关闭旧的连接
        socket?.close()
        socket = null
        _isConnected.value = false
        _isAuthenticated.value = false

        retry(retryCount) { attempt ->
            // 指数退避:每次重试延迟翻倍
            if (attempt > 0) delay(initialDelayMs * (1 shl (attempt - 1)))
            
            socket = client.webSocketSession {
                url(ChatSocketService.Endpoints.ChatSocket.url)
            }

            if (socket?.isActive == true) {
                _isConnected.value = true
                Resource.Success(Unit)
            } else {
                throw IOException("Couldn't establish active connection")
            }
        }
    } catch (e: Exception) {
        _isConnected.value = false
        _isAuthenticated.value = false
        e.printStackTrace()
        Resource.Error("Failed after $retryCount attempts: ${e.localizedMessage ?: "Unknown error"}")
    }
}

方式2:手动实现重试逻辑(无需额外依赖)

override suspend fun initSession(retryCount: Int = 3): Resource<Unit> {
    var attempt = 0
    while (attempt < retryCount) {
        runCatching {
            socket?.close()
            socket = null
            _isConnected.value = false
            _isAuthenticated.value = false

            socket = client.webSocketSession {
                url(ChatSocketService.Endpoints.ChatSocket.url)
            }

            if (socket?.isActive == true) {
                _isConnected.value = true
                return Resource.Success(Unit)
            } else {
                throw IOException("Connection not active")
            }
        }.onFailure {
            attempt++
            if (attempt < retryCount) {
                delay(1000 * attempt) // 每次重试延迟增加1s
            }
        }
    }
    _isConnected.value = false
    _isAuthenticated.value = false
    return Resource.Error("Failed after $retryCount attempts")
}

authenticate()的重试逻辑可以参考同样的方式,或者在Repository层调用时统一处理。


3. 连接状态的持有位置:Service还是Repository?

最佳实践是:让Service层持有连接的真实状态,Repository层负责暴露状态给ViewModel,不持有状态源

核心原因:

  • Service层(ChatSocketServiceImpl)是直接管理WebSocket会话的角色,它最清楚连接的状态变化,状态的维护和更新应该由它负责。
  • Repository层的职责是协调多个数据源(比如本地缓存+网络/WebSocket),它应该将Service的状态Flow转发给ViewModel,而不是自己持有状态副本。

具体实现:

  1. 在Service中将isConnected和isAuthenticated改为不可变的StateFlow(如前面代码所示)
  2. 在ChatRepositoryImpl中暴露这些Flow:
class ChatRepositoryImpl(private val chatSocketService: ChatSocketService) : ChatRepository {
    override fun observeConnectionState(): StateFlow<Boolean> = chatSocketService.isConnected
    override fun observeAuthenticationState(): StateFlow<Boolean> = chatSocketService.isAuthenticated
    // ... 其他Repository方法
}
  1. ViewModel直接观察Repository暴露的Flow,无需关心状态的具体来源。

额外的优化建议

  • 资源清理:在ViewModel的onCleared()中调用Service新增的closeSession()方法,关闭WebSocket连接,避免内存泄漏。
  • 异常处理:在observeMessages()中处理WebSocket断开的情况,比如当incoming流结束时,自动更新isConnected状态。
  • 线程安全:确保对socket、状态Flow的操作都在协程中执行,避免多线程访问的安全问题。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 12:09:54