如何用RxJava2实现BLE通知中图片数据包的正确缓冲?
解决RxJava2中BLE图片分包精准缓冲的问题
我来帮你捋清楚这个问题的核心:你现在的缓冲方案依赖delay+尾包过滤来界定缓冲边界,这种方式在传输中断时完全不可靠——一旦没收到尾包,旧的不完整数据会一直留在缓冲里,和新图片的数据包混在一起。咱们得换一种基于首包/尾包精准分割数据流的思路,用RxJava的window操作符来实现才是正解。
核心思路
用window把原始的BLE字节流分割成一个个独立的子流,每个子流对应从首包到尾包的完整图片数据包。同时处理两种边界情况:
- 正常流程:子流从首包开始,收到尾包时结束
- 异常断包:如果一个子流打开后(收到首包),超时没收到尾包,或者新的首包提前到来,直接终止旧子流,丢弃不完整数据,避免混叠
具体代码实现
首先先定义首包和尾包的判断逻辑(根据你的实际规则调整):
private boolean isFirstPackage(byte[] bytes) { // 替换成你识别首包的实际规则,比如特定标识位 return bytes[0] == 0x01; } private boolean isLastPackage(byte[] bytes) { // 保留你原来的尾包判断逻辑 return bytes[2] < 0; }
然后用window+publish实现精准分包:
// 先把BLE通知流转换成字节数组流 Observable<byte[]> byteStream = notificationObservable() .map(notification -> notification.getBytes()); // 用publish共享流,避免多次订阅原始BLE通知流 disposables.add(byteStream.publish(sharedStream -> sharedStream.window( // 窗口开启条件:收到图片首包 sharedStream.filter(this::isFirstPackage), // 窗口关闭条件:两种情况二选一 window -> Observable.merge( // 情况1:当前窗口收到尾包,正常结束 window.filter(this::isLastPackage), // 情况2:新的首包到来,强制终止旧窗口(处理断包后新图片到来的场景) sharedStream.filter(this::isFirstPackage) ).firstOrError().toObservable() // 额外加超时处理:如果5秒没收到尾包,直接终止窗口(避免内存泄漏) .timeout(5, TimeUnit.SECONDS) ) ) // 将每个窗口的分包合并成完整的图片字节数组 .flatMapSingle(window -> window .collect(ByteArrayOutputStream::new, (baos, bytes) -> { try { baos.write(bytes); } catch (IOException e) { throw new RuntimeException("字节流写入失败", e); } }) .map(ByteArrayOutputStream::toByteArray) ) // 转换成你的自定义图片对象 .map(MyImage::new) .subscribe( imageSubject::onNext, throwable -> { // 这里处理错误:比如超时、IO异常等,可根据需求选择是否通知上层 imageSubject.onError(throwable); }, imageSubject::onComplete ));
方案优势
- 完全避免数据混叠:每个窗口严格对应一张图片的完整数据包,新窗口只会在新首包到来时开启,旧窗口要么正常结束(收到尾包),要么被强制终止(新首包到来/超时)
- 异常场景覆盖:超时处理避免了断包后不完整数据一直占用内存的问题
- 性能更优:用
ByteArrayOutputStream代替ArrayList<Byte>,减少了装箱拆箱的性能损耗
注意事项
- 首包/尾包的判断逻辑一定要准确,这是整个分割逻辑的基础
- 超时时间要根据你的BLE传输速度和图片大小调整(比如图片最大传输时间是3秒,超时设为4秒就足够)
- 如果需要处理断包重试,可以在错误回调里添加对应的重试逻辑,比如重新请求图片传输
内容的提问来源于stack exchange,提问作者Szelk Baltazár
相关产品推荐
相关产品推荐

