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

Android Kotlin使用callbackFlow返回Socket结果报错问题咨询

问题根因

你的代码触发崩溃是四个核心问题导致的:

  • 注册的Socket事件监听没有在流生命周期结束时移除,重复调用会累计无效监听,后续事件触发时向已关闭的通道发送数据,直接触发ClosedSendChannelException崩溃。
  • awaitClose块内逻辑错误,你写的cancel()会在流正常关闭时主动取消协程,该代码块的正确作用是释放流持有的资源(比如移除Socket监听),而非主动取消协程。
  • 异常捕获不全、存在线程竞态:仅捕获了SocketException,回调内的类型强转异常、Socket内部抛出的其他异常都会直接触发崩溃;在回调外定义的response变量存在多线程读写竞态问题。
  • 逻辑不符合一次性调用预期:没有在拿到请求结果后关闭流,Socket后续收到同事件消息时会持续向流发送数据,和你“结果就绪立即返回”的需求不符。
修复后实现
sealed class SocketCallback {
    data class OnSuccess(val data: JSONObject) : SocketCallback()
    data class OnError(val msg: String) : SocketCallback()
}

private fun callSocket(
    eventEmit: String,
    eventOn: String,
    request: JSONObject
) = callbackFlow {
    // 提前定义监听实例,后续用于注册和移除
    val responseListener = Emitter.Listener { args ->
        // 捕获回调内所有异常,避免强转、空指针等崩溃
        val result = runCatching {
            val response = args[0] as JSONObject
            Log.d("SOCKET_RESPONSE", "Event $eventOn received: $response")
            SocketCallback.OnSuccess(response)
        }.getOrElse { e ->
            Log.e("SOCKET_ERROR", "Parse $eventOn response failed", e)
            SocketCallback.OnError("Response parse error: ${e.message}")
        }
        trySend(result)
        close() // 一次性请求拿到结果后关闭流,触发资源释放
    }

    try {
        if (!socket.connected()) {
            val errMsg = "Socket not connected"
            Log.e("SOCKET_ERROR", errMsg)
            trySend(SocketCallback.OnError(errMsg))
            close()
            return@callbackFlow
        }

        Log.d("SOCKET_REQUEST", "Emit event $eventEmit: $request")
        socket.on(eventOn, responseListener)
        socket.emit(eventEmit, request)
    } catch (e: Exception) {
        // 捕获所有可能抛出的异常,避免遗漏
        val errMsg = "Call socket $eventEmit failed: ${e.message}"
        Log.e("SOCKET_ERROR", errMsg, e)
        trySend(SocketCallback.OnError(errMsg))
        close()
    }

    // 流关闭时移除注册的监听,杜绝内存泄漏和无效回调
    awaitClose {
        socket.off(eventOn, responseListener)
    }
}.flowOn(Dispatchers.IO)
使用说明
  • 该实现返回冷流,每次执行collect都会发起一次独立的Socket请求,不会残留重复监听。
  • 如需添加请求超时,可通过流操作符实现,示例:
lifecycleScope.launch {
    callSocket("user_login", "login_result", loginRequest)
        .timeout(10, TimeUnit.SECONDS)
        .catch { e ->
            emit(SocketCallback.OnError("Request timeout: ${e.message}"))
        }
        .collect { res ->
            when(res) {
                is SocketCallback.OnSuccess -> { /* 处理成功结果 */ }
                is SocketCallback.OnError -> { /* 处理错误 */ }
            }
        }
}
  • 如果你需要持续监听某个Socket事件做长连接推送,删掉拿到结果后close()的调用即可,注意在页面销毁、不需要监听时主动取消协程,会自动触发awaitClose移除监听。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 15:24:20