基于RxJava实现网络查询持久化及后续查询必要性校验
这个场景我太熟悉了——既要确保数据持久化的insertAll()(返回Completable)执行,又不能丢掉第一次API请求的结果用来判断是否需要后续查询,确实有点绕。先帮你拆解下原来代码的问题:
你最初的flatMap里只是调用了insertAll(it)但没把它纳入数据流,因为flatMap需要你返回一个Observable/Flowable,而你最后返回的是shouldGetMore(it)的结果,所以insertAll()返回的Completable根本没被订阅,自然不会执行。换成flatMapCompletable虽然能执行insertAll(),但会把上游的Observable转换成Completable,直接丢掉了API返回的数据,后续的判断也就无从谈起。
给你两个不用修改insertAll()返回值的优雅方案:
方案一:用andThen在Completable完成后传递原数据
这是最贴合你需求的方案,利用Completable.andThen()方法,等待持久化操作完成后,把原API数据重新发射到下游,这样既保证了insertAll()执行,又保留了后续判断需要的数据:
return remote(amount = 2) .subscribeOn(Schedulers.io()) .flatMap { data -> // 先执行持久化,完成后把原数据传递给下游 insertAll(data) .andThen(Observable.just(data)) } .flatMap { data -> // 用原数据判断是否需要发起更多查询 if (shouldGetMore(data)) remote(amount = 3) else Observable.just(emptyList()) } .flatMapCompletable { insertAll(it) } .observeOn(AndroidSchedulers.mainThread())
andThen()会等待前面的Completable执行完成(包括成功或失败),然后订阅并发射后面Observable的数据,完美解决了"执行无返回操作+保留上游数据"的矛盾。
方案二:如果持久化是可容忍失败的副作用(不推荐用于核心逻辑)
如果你能接受insertAll()失败不影响整个数据流的执行(比如只是日志类操作),可以用doOnNext来执行副作用,但注意:这个方案不适合核心持久化逻辑,因为doOnNext里的错误不会中断数据流,而且它是在当前线程执行(如果需要切换线程要额外处理):
return remote(amount = 2) .subscribeOn(Schedulers.io()) .doOnNext { insertAll(it).subscribe() } // 手动订阅Completable .flatMap { data -> if (shouldGetMore(data)) remote(amount = 3) else Observable.just(emptyList()) } .flatMapCompletable { insertAll(it) } .observeOn(AndroidSchedulers.mainThread())
这个方案的弊端很明显:手动订阅容易遗漏错误处理,而且如果insertAll()执行慢,下游的判断可能先执行,不符合你"先持久化再判断"的顺序,所以优先选方案一。
总结一下,方案一既不用修改现有insertAll()的返回值,又能严格保证执行顺序和数据传递,是最优雅的解决方式。
内容的提问来源于stack exchange,提问作者nmu

