RxAndroidBle实现BLE加密流程优化求助:避免重复连接
重构RxAndroidBle流程:单一数据流+兼容现有客户端接口
嘿,我明白你现在的困扰——重复建立BLE连接不仅低效,而且完全没发挥出RxJava响应式的优势。咱们可以通过几个RxJava的核心技巧来重构,既保持单一数据流,又兼容你已有的客户端代码,不用改第三方应用的逻辑。
核心思路
- 用Subject桥接命令式回调:因为你的客户端是通过
onKeyNeeded()回调+setKey()方法传入密钥,这是命令式的调用,我们需要用PublishSubject把它转换成Observable流,让RxJava能感知到密钥的到来。 - 保持BLE连接不中断:建立一次连接后,在同一个
RxBleConnection实例上完成所有操作(读加密特征、写密钥、读后续特征),避免重复连接的开销。 - 用
flatMap做分支处理:根据加密特征的读取结果,动态切换数据流分支——已加密直接走读特征流程,未加密则等待密钥写入后再走相同流程。 - 统一管理订阅生命周期:用一个Disposable来管理整个数据流的订阅,确保出错或完成时正确释放连接。
重构后的代码实现
import io.reactivex.rxjava3.core.Observable; import io.reactivex.rxjava3.core.Single; import io.reactivex.rxjava3.disposables.Disposable; import io.reactivex.rxjava3.subjects.PublishSubject; public class BleManager { private final RxBleDevice bleDevice; private final KeyNeededListener listener; private final PublishSubject<byte[]> keySubject = PublishSubject.create(); private Disposable connectionDisposable; // 保留你已有的接口 public interface KeyNeededListener { void onKeyNeeded(); } public BleManager(RxBleDevice bleDevice, KeyNeededListener listener) { this.bleDevice = bleDevice; this.listener = listener; } public void connect() { // 取消之前的订阅,避免泄漏 if (connectionDisposable != null && !connectionDisposable.isDisposed()) { connectionDisposable.dispose(); } connectionDisposable = bleDevice.establishConnection(false) .flatMapSingle(connection -> // 第一步:读取加密特征 connection.readCharacteristic(ENCRYPT_CHARACTERISTIC_UUID) .flatMap(encryptBytes -> { if (isEncrypted(encryptBytes)) { // 分支1:已加密,直接读取后续特征 return readAllDataCharacteristics(connection); } else { // 分支2:未加密,通知客户端需要密钥,然后等待密钥传入 listener.onKeyNeeded(); return keySubject.take(1) // 只取一次传入的密钥 .flatMapSingle(key -> // 写入密钥后,再读取后续特征 connection.writeCharacteristic(ENCRYPT_CHARACTERISTIC_UUID, key) .flatMap(writeResult -> readAllDataCharacteristics(connection)) ); } }) ) .subscribe(this::processData, this::onError, () -> { // 数据流完成时的清理 if (!keySubject.isDisposed()) { keySubject.dispose(); } }); } // 保留你已有的setKey接口,把密钥发射到Subject中 public void setKey(byte[] key) { if (!keySubject.isDisposed()) { keySubject.onNext(key); } } // 提取重复的读特征逻辑,复用代码 private Single<DataModel> readAllDataCharacteristics(RxBleConnection connection) { return Single.zip( connection.readCharacteristic(A_CHARACTERISTIC_UUID), connection.readCharacteristic(B_CHARACTERISTIC_UUID), connection.readCharacteristic(C_CHARACTERISTIC_UUID), DataModel::new ); } // 你的原有方法:判断是否加密 private boolean isEncrypted(byte[] bytes) { // 这里是你的原有逻辑 return false; } // 你的原有方法:处理数据 private void processData(DataModel dataModel) { // 这里是你的原有逻辑 } // 你的原有方法:错误处理 private void onError(Throwable throwable) { // 这里是你的原有逻辑 if (!keySubject.isDisposed()) { keySubject.dispose(); } } }
关键部分解释
PublishSubject<byte[]> keySubject:作为命令式setKey()和响应式数据流之间的桥梁,当客户端调用setKey()时,密钥会被发射到这个Subject中,数据流就能感知到并继续执行。- 保持单一连接:整个流程只调用一次
establishConnection(false),后续所有操作都复用同一个RxBleConnection实例,避免了重复连接的耗时和资源浪费。 flatMap分支处理:读取加密特征后,根据结果动态选择后续的数据流——已加密直接读数据,未加密则等待密钥写入后再读数据,完全符合RxJava的链式响应式设计。take(1):确保我们只使用客户端传入的第一个密钥,避免后续误调用setKey()影响流程。- 生命周期管理:在
connect()开始时取消之前的订阅,在数据流完成或出错时清理keySubject,避免内存泄漏。
这样重构后,你既保留了原有的客户端接口(不用修改第三方应用的代码),又实现了单一数据流的响应式设计,完全消除了重复连接的问题。
内容的提问来源于stack exchange,提问作者N. Smith
相关产品推荐
相关产品推荐

