RxJava 1.0:如何阻塞至嵌套订阅中所有Google调用完成?
看起来你遇到的是RxJava中嵌套订阅导致主流无法等待异步操作完成的典型问题——你用toBlocking只阻塞了外层流,但内层的Google API调用是在独立的订阅里异步执行的,外层流的doOnCompleted会提前触发,这时候缓存还没被填充。
问题根源
你原来的代码里,group.subscribe(...)是在创建新的独立订阅,外层的cartesian流根本不会等待这些内层订阅完成。哪怕用了toBlocking,也只是等外层流(分组完成)结束,而内层的Google调用还在后台跑,自然没法在doOnCompleted里拿到完整的缓存。
解决方案:避免嵌套订阅,合并流
解决的核心是把内层的Google调用流合并到外层主流中,让整个流的完成信号等待所有异步操作结束。用flatMap(或concatMap如果需要顺序执行)代替嵌套的subscribe,这样所有操作都在同一个流链里,toBlocking就能正确阻塞到所有任务完成。
修改后的代码示例
假设你原来的内层逻辑是对每个分组里的CartesianProduct调用Google API,然后把结果存入缓存,修改后的代码大概是这样:
cartesian // 先过滤掉缓存已有的项 .filter(cp -> !cache.containsKey(cp.fromTo)) // 按from字段分组 .groupBy(cp -> cp.from) // 将每个分组的处理流合并到主流 .flatMap(group -> group.flatMap(cp -> // 这里是你的Google API调用,需返回RxJava流(Observable/Flowable等) googleApi.getGeoLocation(cp) // API调用完成后更新缓存 .doOnNext(result -> cache.put(cp.fromTo, result)) ) ) // 所有Google调用+缓存更新完成后才会触发这个回调 .doOnCompleted(() -> { System.out.println("DOC#1 # ELEMENTS IN CACHE " + cache.size()); // 这里可以安全访问缓存,所有操作都已完成 }) // 阻塞直到整个流的所有操作执行完毕 .toBlocking() .subscribe();
为什么这样有效?
flatMap会将每个分组的内层流(Google API调用)合并到外层主流中,主流的onCompleted只有在所有内层流都执行完毕后才会触发。toBlocking现在阻塞的是整个合并后的流,所以会等待所有Google API调用和缓存更新操作完成后才继续执行。- 如果需要保证分组内的调用顺序(比如按顺序处理每个
from下的请求),可以把flatMap换成concatMap,它会按顺序处理每个分组的流。
额外注意事项
- 如果你的Google API调用是同步方法,可以用
Observable.fromCallable()把它包装成异步流,比如:Observable.fromCallable(() -> googleApi.syncCall(cp))。 - 别忘了添加错误处理,比如在流链中加入
onErrorResumeNext或onErrorReturn,避免单个API调用失败导致整个流中断。
内容的提问来源于stack exchange,提问作者user2849678
相关产品推荐
相关产品推荐

