RxJava:如何展开嵌套Observable并实现定时服务请求?
嘿,我来帮你搞定这个RxJava的定时请求问题~
你现在的情况是:用Observable.interval每分钟触发一次请求,但因为restService.getData()返回的是Single<MyData>,所以最终得到的是Observable<Single<MyData>>——这确实不是你想要的,你需要的是直接拿到每次请求返回的MyData来处理。
其实不用绕到Single.fromObservable,RxJava有更简洁的操作符来解决这个问题,我给你两种常用的方案:
1. 并发执行定时请求(默认推荐)
用flatMapSingle操作符,它专门用来把Observable的每个事件转换成Single,然后自动订阅这些Single,把它们的发射结果合并成一个Observable<MyData>。这样你就能直接订阅这个Observable,处理每次请求返回的数据了:
Observable.interval(1, TimeUnit.MINUTES, Schedulers.io()) .flatMapSingle(tick -> restService.getData()) .subscribe( myData -> { // 这里处理每次拿到的MyData System.out.println("获取到数据:" + myData.toString()); }, error -> { // 处理全局流的错误(比如interval调度失败) System.err.println("定时任务出错:" + error.getMessage()); } );
注意错误处理:
如果某次请求失败不想终止整个定时流,可以在flatMapSingle里给Single添加错误兜底:
Observable.interval(1, TimeUnit.MINUTES, Schedulers.io()) .flatMapSingle(tick -> restService.getData() .onErrorResumeNext(error -> { // 处理单个请求的错误,比如打日志 System.err.println("本次请求失败:" + error.getMessage()); // 返回empty表示忽略这次错误,继续下一次定时 return Single.empty(); // 如果想把错误抛出去,就用return Single.error(error),但这样会终止整个流 }) ) .subscribe( myData -> { /* 处理有效数据 */ }, error -> { /* 全局流错误处理 */ } );
2. 顺序执行定时请求(等上一次完成再发起下一次)
如果你的业务要求必须等上一次请求完成后,再等1分钟发起下一次(避免并发请求),可以用concatMapSingle替代flatMapSingle,同时把interval换成timer(确保第一次请求立即执行,之后每隔1分钟执行一次):
Observable.timer(0, 1, TimeUnit.MINUTES, Schedulers.io()) .concatMapSingle(tick -> restService.getData()) .subscribe( myData -> { /* 处理数据 */ }, error -> { /* 错误处理 */ } );
关于你原来的Single.fromObservable思路
其实这个思路也能走通,但只适合你只需要第一次请求的数据的场景(因为Single只会发射一次),代码会更啰嗦:
Single.fromObservable( Observable.interval(1, TimeUnit.MINUTES, Schedulers.io()) .flatMapSingle(tick -> restService.getData()) .take(1) // 只取第一次数据 ) .subscribe( myData -> { /* 处理第一次数据 */ }, error -> { /* 错误处理 */ } );
但显然这不符合你“定时请求”的核心需求,所以还是前面两种方案更合适~
内容的提问来源于stack exchange,提问作者Dmitry Zhgun
相关产品推荐
相关产品推荐

