展平嵌套Observable:实现每秒执行函数返回Observable<A>
解决RxJava嵌套Observable扁平化问题
我来帮你搞定这个需求——实现一个每秒调用返回Observable<A>的函数,同时把嵌套的Observable<Observable<A>>转换成单一的Observable<A>流。下面是具体的思路和代码示例:
核心实现思路
RxJava里专门有操作符用来处理嵌套Observable的扁平化问题,再配合定时触发的操作符就能完美实现你的需求:
interval():实现每秒触发一次调用的定时逻辑concatMap/flatMap:把每个调用返回的Observable<A>扁平化成统一的流(concatMap保证顺序,flatMap支持并行)takeWhile():替代你示例里的while(condition),用来控制流的终止时机
代码示例(适配返回Observable的目标函数)
假设你的calcA()现在返回Observable<A>(符合你问题描述的需求),改造后的getAs()可以这样写:
// 目标函数:返回Observable<A> private Observable<A> calcA() { // 这里替换成你的实际业务逻辑 return Observable.just(new A(...)); } public Observable<A> getAs() { return Observable.interval(0, 1, TimeUnit.SECONDS) // 控制循环终止条件:condition为true时继续发射,false则终止流 .takeWhile(tick -> { boolean condition = ...; // 你的终止判断逻辑 return condition; }) // 扁平化嵌套Observable,用concatMap保证调用顺序和结果顺序一致 .concatMap(tick -> calcA()) .doOnComplete(() -> System.out.println("数据流已终止")); }
关键细节说明
interval(0, 1, TimeUnit.SECONDS):第一个参数0表示订阅后立即执行第一次调用,之后每隔1秒触发一次;如果想等待1秒再开始,把第一个参数改成1即可。concatMapvsflatMap:如果calcA()是耗时操作且不关心结果顺序,用flatMap可以提升并行处理效率;如果需要严格按照每秒调用的顺序输出结果,concatMap是更稳妥的选择。- 如果calcA()返回的是A而非Observable(和你示例代码一致),那更简单,把
concatMap换成map就行:
private A calcA() { return new A(...); } public Observable<A> getAs() { return Observable.interval(0, 1, TimeUnit.SECONDS) .takeWhile(tick -> { boolean condition = ...; return condition; }) .map(tick -> calcA()); }
内容的提问来源于stack exchange,提问作者godzsa
相关产品推荐
相关产品推荐

