You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJava Zip操作符并行失效:任务均运行于主线程的问题排查

问题原因

  1. 原始代码中,getData()和getData3()方法内的阻塞逻辑(Thread.sleep、线程打印)是同步执行的——调用repository.getData()时这些代码会立即运行,而此时调用线程是主线程,所以逻辑直接跑在主线程,还会串行执行(先等getData的10秒sleep结束,才会执行getData3)。
  2. subscribeOn(Schedulers.io())只能控制Observable订阅后的数据发射流程,没法改变调用repository.getData()这个方法本身的线程。
  3. 后续的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 00:20:31