基于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,而不是自己持有状态副本。
具体实现:
- 在Service中将
isConnected和isAuthenticated改为不可变的StateFlow(如前面代码所示) - 在
ChatRepositoryImpl中暴露这些Flow:
class ChatRepositoryImpl(private val chatSocketService: ChatSocketService) : ChatRepository { override fun observeConnectionState(): StateFlow<Boolean> = chatSocketService.isConnected override fun observeAuthenticationState(): StateFlow<Boolean> = chatSocketService.isAuthenticated // ... 其他Repository方法 }
- ViewModel直接观察Repository暴露的Flow,无需关心状态的具体来源。
额外的优化建议
- 资源清理:在ViewModel的
onCleared()中调用Service新增的closeSession()方法,关闭WebSocket连接,避免内存泄漏。 - 异常处理:在
observeMessages()中处理WebSocket断开的情况,比如当incoming流结束时,自动更新isConnected状态。 - 线程安全:确保对
socket、状态Flow的操作都在协程中执行,避免多线程访问的安全问题。
内容来源于stack exchange
相关产品推荐
相关产品推荐

