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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 06:47:14