RxJava Zip操作符并行失效:任务均运行于主线程的问题排查
问题原因
- 原始代码中,
getData()和getData3()方法内的阻塞逻辑(Thread.sleep、线程打印)是同步执行的——调用repository.getData()时这些代码会立即运行,而此时调用线程是主线程,所以逻辑直接跑在主线程,还会串行执行(先等getData的10秒sleep结束,才会执行getData3)。 subscribeOn(Schedulers.io())只能控制Observable订阅后的数据发射流程,没法改变调用repository.getData()这个方法本身的线程。- 后续的defer修改无效,是因为阻塞逻辑写在了
Single.defer()调用之前,仍然在主线程同步执行,defer只推迟了Single.just(1)的创建,根本没覆盖耗时逻辑。
解决方案
核心是把所有耗时逻辑和Observable/Single的创建逻辑包裹在能推迟执行的操作符中,让这些逻辑在subscribeOn指定的io线程上运行,而非调用方法的主线程。推荐两种实现方式:
方式1:使用defer操作符
class FakeRepositoryImpl @Inject constructor() : FakeRepository { override fun getData(): Observable<Int> { return Observable.defer { println("test2k200: getData " + Thread.currentThread().name) Thread.sleep(10000) Observable.just(1) } } override fun getData3(): Observable<Int> { return Observable.defer { println("test2k200: getData3 " + Thread.currentThread().name) Observable.just(3) } } }
方式2:使用fromCallable操作符(更简洁)
fromCallable专门用于把同步逻辑包装成Observable,还能自动处理异常,比defer+just更直观:
class FakeRepositoryImpl @Inject constructor() : FakeRepository { override fun getData(): Observable<Int> { return Observable.fromCallable { println("test2k200: getData " + Thread.currentThread().name) Thread.sleep(10000) 1 // 直接返回结果,自动包装为Observable } } override fun getData3(): Observable<Int> { return Observable.fromCallable { println("test2k200: getData3 " + Thread.currentThread().name) 3 } } }
订阅代码无需修改(保持原zip逻辑即可)
fun getDataSync() { compositeDisposable.add( Observable.zip( repository.getData().subscribeOn(Schedulers.io()), repository.getData3().subscribeOn(Schedulers.io()), { a, b -> a + b } ) .observeOn(AndroidSchedulers.mainThread()) .subscribe( { _data.value = it }, { _data.value = 0 } ) ) }
效果验证
修改后,控制台打印的线程会是RxCachedThreadScheduler-x(属于io线程池),且两个方法会几乎同时启动执行,不会出现先等10秒再执行第二个的情况,真正实现并行。
内容的提问来源于stack exchange,提问作者discCard
相关产品推荐
相关产品推荐

