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

