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

使用callbackFlow将WebRTC PeerConnection.Observer转为Kotlin Flows

用Kotlin Flow优雅处理WebRTC PeerConnection.Observer多回调事件

针对WebRTC PeerConnection.Observer的多回调场景,我们可以通过密封类统一事件类型 + 单个callbackFlow分发所有事件的方式实现简洁高效的Flow封装,既保留Flow的响应式特性,又能清晰处理不同类型的WebRTC回调。

步骤1:定义密封类封装所有WebRTC事件

先把PeerConnection.Observer中需要监听的回调事件封装成密封类,每个事件对应一个子类,携带回调的参数:

sealed class WebRtcEvent {
    data class IceCandidate(val candidate: IceCandidate) : WebRtcEvent()
    data class IceConnectionStateChanged(val state: PeerConnection.IceConnectionState) : WebRtcEvent()
    data class SignalingStateChanged(val state: PeerConnection.SignalingState) : WebRtcEvent()
    data class OnAddStream(val stream: MediaStream) : WebRtcEvent()
    data class OnRemoveStream(val stream: MediaStream) : WebRtcEvent()
    // 根据需求添加其他需要的回调事件,比如OnDataChannel、OnIceGatheringStateChanged等
}

步骤2:实现PeerConnection.Observer并绑定callbackFlow

创建一个实现PeerConnection.Observer的类,内部持有callbackFlow的SendChannel,在每个回调方法中发送对应的密封类事件:

class FlowPeerConnectionObserver(
    private val sendChannel: SendChannel<WebRtcEvent>
) : PeerConnection.Observer {

    // 处理ICE候选回调
    override fun onIceCandidate(candidate: IceCandidate) {
        sendChannel.trySend(WebRtcEvent.IceCandidate(candidate))
    }

    // 处理ICE连接状态变化
    override fun onIceConnectionStateChanged(newState: PeerConnection.IceConnectionState) {
        sendChannel.trySend(WebRtcEvent.IceConnectionStateChanged(newState))
    }

    // 处理信令状态变化
    override fun onSignalingStateChanged(newState: PeerConnection.SignalingState) {
        sendChannel.trySend(WebRtcEvent.SignalingStateChanged(newState))
    }

    // 处理添加媒体流
    override fun onAddStream(stream: MediaStream) {
        sendChannel.trySend(WebRtcEvent.OnAddStream(stream))
    }

    // 处理移除媒体流
    override fun onRemoveStream(stream: MediaStream) {
        sendChannel.trySend(WebRtcEvent.OnRemoveStream(stream))
    }

    // 其他不需要的回调直接使用WebRTC提供的默认实现即可
}

步骤3:创建并暴露Flow供上层订阅

通过callbackFlow创建Flow,初始化FlowPeerConnectionObserver并绑定到PeerConnection,同时处理Flow的生命周期:

fun createPeerConnectionFlow(peerConnection: PeerConnection): Flow<WebRtcEvent> = callbackFlow {
    val observer = FlowPeerConnectionObserver(channel)
    
    // 将Observer绑定到PeerConnection
    peerConnection.registerObserver(observer)

    // 当Flow取消时,解绑Observer并关闭通道,避免内存泄漏
    awaitClose {
        peerConnection.unregisterObserver(observer)
        channel.close()
    }
}

步骤4:上层订阅并处理事件

在业务代码中订阅Flow,通过when分支处理不同类型的WebRTC事件:

// 假设已初始化peerConnection
lifecycleScope.launch {
    createPeerConnectionFlow(peerConnection)
        .flowOn(Dispatchers.IO) // WebRTC回调运行在内部线程池,切换到IO线程处理
        .collect { event ->
            when (event) {
                is WebRtcEvent.IceCandidate -> {
                    // 处理ICE候选,比如发送给远端
                }
                is WebRtcEvent.IceConnectionStateChanged -> {
                    // 更新UI显示连接状态
                }
                is WebRtcEvent.OnAddStream -> {
                    // 渲染远端媒体流
                }
                // 处理其他事件...
            }
        }
}

关键注意事项

  • 线程切换:WebRTC的回调方法运行在内部线程池,订阅Flow时建议用flowOn(Dispatchers.IO),再通过lifecycleScope切换回主线程处理UI操作。
  • 事件发送安全性:使用trySend而非send,避免因Flow缓冲区已满导致阻塞WebRTC内部线程。
  • 生命周期管理:通过awaitClose在Flow取消时解绑Observer并关闭通道,防止内存泄漏。
  • 按需实现回调:只实现PeerConnection.Observer中需要的方法,其他方法使用默认实现,减少冗余代码。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 14:55:24