Android RxJava中Flowable.fromIterable不结束 无法触发subscribe问题求助
问题根因排查
核心原因:上游流未发送onComplete事件,toList()永远等待
toList()操作符必须等到上游流所有元素发射完成、且收到onComplete回调后,才会组装成列表发射给下游。你的流程卡住90%以上是因为getAttachedDocs()返回的流没有发送onComplete事件:
如果你的getAttachedDocs()是Room这类ORM的可观察查询(返回Flowable<List<DocumentoAdjunto>>),这类流默认是常驻的,只要数据库对应数据发生变更就会发射新的列表,永远不会主动发送onComplete,导致toList()一直处于等待状态,下游的subscribe自然不会触发。
其他隐藏问题
- 删除操作和主流程完全脱节:
deleteDoc()内部单独做了订阅,和上传主流程没有关联,会出现主流程已经执行完成、删除操作还在后台执行的时序问题,删除报错也不会被主流程捕获。 - 单文件上传错误会打断整个流:
sendToServer()如果抛出异常没有捕获,会直接触发全局onError,不会继续处理剩余文件。 - 部分场景下
sendToServer()返回的Flowable可能不发送onComplete:如果uploadFile()是常驻Observable,转成Flowable后也会出现和getAttachedDocs()一样的不结束问题。
修复方案
1. 修复上游流不结束问题
给getAttachedDocs()返回的Flowable加take(1),只取第一次查询结果,拿到后主动结束流:
public @NonNull Single<List<Boolean>> sendMediaAttachement() { return mRepository.getAttachedDocs(getMantenValue().getID()) .take(1) // 仅取一次数据,之后主动发送onComplete .concatMap(docs -> Flowable.fromIterable(docs)) .flatMap(doc -> mRepository.sendToServer(doc, getFileFromAttachedDoc(doc)) // 单个文件上传失败不打断整个流程 .onErrorReturn(e -> { Timber.e(e); return WSResult.failure(); }) .flatMap(wsResult -> { if (wsResult.isSuccessful()) { // 将删除操作合并到主流程,等待删除完成再往下走 return deleteDoc(doc).map(ignored -> true); } return Flowable.just(false); }) ) .toList(); }
2. 改造deleteDoc,和主流程关联
删除内部订阅逻辑,将删除结果返回给主流程:
public Single<Boolean> deleteDoc(DocumentoAdjunto doc) { return mRepository.delete(doc) .doOnSuccess(deleted -> { if (deleted) FileUtils.deleteFile(getFileFromAttachedDoc(doc)); }) .doOnError(t -> getToastMessageInteger().setValue(R.string.file_not_deleted)) // 即使删除失败也不打断主流程,可根据业务需求调整 .onErrorReturnItem(false) .toFlowable() .firstOrError(); }
3. 确保sendToServer的流一定会结束
给返回的Flowable加take(1),避免上传流常驻:
public @NonNull Flowable<WSResult> sendToServer(DocumentoAdjunto doc, File file) { return mRemoteDataSource.uploadFile(doc, file) .map(this::parse) .toFlowable() .take(1) // 确保发射一次结果后主动结束 .doOnError(Timber::e) .subscribeOn(Schedulers.io()); }
调试建议
如果修改后还是有问题,可以在每个操作符后加日志定位流的执行状态:
.concatMap(docs -> Flowable.fromIterable(docs)) .doOnNext(doc -> Log.d("RX_DEBUG", "开始处理文档:" + doc.getId())) .doOnComplete(() -> Log.d("RX_DEBUG", "所有文档遍历完成")) .flatMap(// 原有逻辑) .doOnComplete(() -> Log.d("RX_DEBUG", "所有上传删除操作完成")) .toList() .doOnSuccess(list -> Log.d("RX_DEBUG", "结果列表长度:" + list.size()))
内容的提问来源于stack exchange,提问作者mantc_sdr
相关产品推荐
相关产品推荐

