RxJava:为combineLatest内各Observable设置不同调度器及验证方法
为combineLatest中不同Observable指定独立调度器的实现与验证
可行性结论
完全可行。RxJava中每个Observable可以通过subscribeOn()独立指定自身发射数据的执行线程,彼此之间互不干扰,下游额外添加的subscribeOn()在这里无意义,因为上游每个Observable已经明确了自身的执行线程。
实现代码
直接给每个Observable单独添加subscribeOn()指定对应调度器,同时保留observeOn()指定主线程处理合并结果即可:
Observables.combineLatest( // 给observable1指定computation调度器 observable1.subscribeOn(schedulers.computation()), // 给observable2指定io调度器 observable2.subscribeOn(schedulers.io()) ) // 合并后的数据在主线程处理 .observeOn(schedulers.mainThread()) .subscribeBy( onNext = { (data1, data2) -> // 处理合并后的业务逻辑 }, onError = { error -> // 错误处理逻辑 } )
验证方法:打印线程日志
通过打印每个Observable发射数据时的线程名称,以及合并后处理数据的线程名称,就能直观验证调度器是否生效:
- 改造Observable,添加线程日志
// 改造observable1,打印发射线程 val observable1 = Observable.create<String> { emitter -> println("Observable1 发射线程: ${Thread.currentThread().name}") emitter.onNext("测试数据1") emitter.onComplete() }.subscribeOn(schedulers.computation()) // 改造observable2,打印发射线程 val observable2 = Observable.create<String> { emitter -> println("Observable2 发射线程: ${Thread.currentThread().name}") emitter.onNext("测试数据2") emitter.onComplete() }.subscribeOn(schedulers.io())
- 在订阅回调中打印处理线程
.subscribeBy( onNext = { (data1, data2) -> println("合并后处理线程: ${Thread.currentThread().name}") // 业务逻辑 } )
预期日志输出
- Observable1的线程名称会包含
RxComputationThreadPool(对应computation调度器) - Observable2的线程名称会包含
RxCachedThreadScheduler(对应io调度器) - 合并后处理线程会显示主线程名称(比如Android环境下的
main线程)
内容的提问来源于stack exchange,提问作者Rainmaker
相关产品推荐
相关产品推荐

