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

RxAndroidBle实现BLE加密流程优化求助:避免重复连接

重构RxAndroidBle流程:单一数据流+兼容现有客户端接口

嘿,我明白你现在的困扰——重复建立BLE连接不仅低效,而且完全没发挥出RxJava响应式的优势。咱们可以通过几个RxJava的核心技巧来重构,既保持单一数据流,又兼容你已有的客户端代码,不用改第三方应用的逻辑。

核心思路

  1. 用Subject桥接命令式回调:因为你的客户端是通过onKeyNeeded()回调+setKey()方法传入密钥,这是命令式的调用,我们需要用PublishSubject把它转换成Observable流,让RxJava能感知到密钥的到来。
  2. 保持BLE连接不中断:建立一次连接后,在同一个RxBleConnection实例上完成所有操作(读加密特征、写密钥、读后续特征),避免重复连接的开销。
  3. 用flatMap做分支处理:根据加密特征的读取结果,动态切换数据流分支——已加密直接走读特征流程,未加密则等待密钥写入后再走相同流程。
  4. 统一管理订阅生命周期:用一个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 06:46:49