如何将BleDataSource的Lambda回调转换为Flow返回?
实现方案
1. 调整BleOperationType和Connect数据类
移除原来的回调参数,改为持有Flow的发射器(Emitter),让后续回调能直接向Flow发送状态:
abstract class BleOperationType { abstract val emitter: Emitter<Resource<BleOperationResult>> } data class Connect( val device: BluetoothDevice, override val emitter: Emitter<Resource<BleOperationResult>> ) : BleOperationType()
2. 修改BleDataSource的performConnect方法,返回Flow
用callbackFlow把回调逻辑转为Flow,这是Kotlin协程处理回调转Flow的标准方式:
// 替换原有performConnect方法 fun performConnect(device: BluetoothDevice): Flow<Resource<BleOperationResult>> = callbackFlow { val connectOperation = Connect(device, this) enqueueOperation(connectOperation) // Flow取消时清理资源:终止当前BLE连接 awaitClose { val pendingOp = pendingOperation if (pendingOp is Connect && pendingOp.device == device) { pendingOp.device.bluetoothGatt?.disconnect() signalEndOfOperation() } } }
3. 调整BleDataSource内的状态发送逻辑
把所有原来调用operation.result()的地方,改为通过emitter向Flow发送状态,同时处理Flow的关闭:
3.1 修改doNextOperation中的初始Loading状态发送
private fun doNextOperation() { // ... 原有代码 ... if ( operation is Connect ) { with(operation) { // 替换原result调用,用trySend避免Flow已取消时抛出异常 emitter.trySend(Resource.Loading(message = "正在连接${device.name}")) bluetoothGatt = if ( Build.VERSION.SDK_INT < Build.VERSION_CODES.M ) { device.connectGatt(context, false, gattCallback) } else { device.connectGatt(context, false, gattCallback, BluetoothDevice.TRANSPORT_LE) } } } }
3.2 修改onConnectionStateChange回调
override fun onConnectionStateChange(gatt: BluetoothGatt, status: Int, newState: Int) { val deviceAddress = gatt.device.address val operation = pendingOperation var res: Resource<BleOperationResult> = Resource.Error(errorMessage = "未知错误!") if (status == BluetoothGatt.GATT_SUCCESS) { if (newState == BluetoothProfile.STATE_CONNECTED) { res = Resource.Loading(message = "正在发现服务") gatt.discoverServices() } else if (newState == BluetoothProfile.STATE_DISCONNECTED) { res = Resource.Error(errorMessage = "意外断开连接") } } else { res = Resource.Error(errorMessage = "设备$deviceAddress 出现错误:$status!") } if (operation is Connect) { operation.emitter.trySend(res) } if (res is Resource.Error) { if (operation is Connect) { operation.emitter.close() signalEndOfOperation() } } }
3.3 修改onServicesDiscovered回调
override fun onServicesDiscovered(gatt: BluetoothGatt?, status: Int) { val operation = pendingOperation var res: Resource<BleOperationResult> = Resource.Error(errorMessage = "未知错误!") if (status == BluetoothGatt.GATT_SUCCESS) { res = Resource.Success(data = BleOperationResult.ConnectionResult(profile)) } else { res = Resource.Error(errorMessage = "发现服务失败...") } if (operation is Connect) { operation.emitter.trySend(res) operation.emitter.close() } if (pendingOperation is Connect) { signalEndOfOperation() } }
4. 仓库层实现返回Flow的connect方法
直接调用修改后的performConnect即可:
override fun connect(device: BluetoothDevice): Flow<Resource<BleOperationResult>> { return handler.performConnect(device) }
关键说明
callbackFlow是回调转Flow的最优方案,自带协程上下文管理,awaitClose能确保Flow取消时及时清理BLE连接资源。- 使用
trySend而非send,避免Flow已取消(如页面销毁)时抛出异常。 - 每个连接操作对应独立Flow,流程结束(成功/失败)时手动关闭Flow,防止内存泄漏。
内容的提问来源于stack exchange,提问作者Metropol_Tilkisi
相关产品推荐
相关产品推荐

