订阅多个RxJava Observable时部分订阅者onNext()未触发?
你的问题看起来是典型的竞态条件引发的日志/事件异常,但和你最初理解的“订阅者未就绪”略有差异,我来拆解一下:
首先,你的代码逻辑本身是没问题的——循环创建6个独立Observable,每个都通过subscribeOn把发射逻辑放到计算线程,再通过observeOn切换回主线程订阅。但出现日志丢失的核心原因可能有两个:
日志框架的合并/丢弃行为
你所有Observable发射的都是相同的字符串"somestring",对应的日志输出也是完全一致的"testIt onNext"。很多日志框架(比如常见的Timber、自定义Logger)会对短时间内重复的日志做合并优化,避免刷屏。当6个日志事件几乎同时涌入主线程时,框架可能只保留1-2条重复输出;而你在onNext中添加延迟后,日志事件有了时间间隔,框架就会逐个输出。计算线程池的调度特性引发的极端竞态
Schedulers.computation()的线程数量等于CPU核心数(比如4核设备只有4个线程),当你快速提交6个任务时,后2个会进入排队状态。每个任务的发射逻辑极短(立即onNext+onComplete),可能在某些极端情况下,任务执行完成后,主线程的观察者还没完全完成事件绑定——不过这种情况在RxJava的标准实现里非常罕见,优先级低于日志框架的可能性。
解决建议
1. 先验证日志是否真的丢失(而非被合并)
修改日志内容,加入循环索引i来区分每个Observable的事件:
.subscribe({ Logger.i("testIt onNext from Observable $i") }, { Logger.i("testIt onError") })
这样你就能明确看到每个Observable的onNext是否真的被触发,而不是被日志框架合并成了相同输出。
2. 调整线程池类型,减少任务排队
计算线程池是为CPU密集型任务设计的,线程数量有限;换成IO线程池(专为短任务设计,线程数量更大)可以减少任务排队的概率:
.subscribeOn(Schedulers.io())
3. 给Observable添加微小延迟,规避极端竞态
如果确实是RxJava内部的竞态导致事件丢失,可以用delay操作符给发射逻辑加个微小延迟,给订阅流程足够的绑定时间:
Observable.create<String>(object : ObservableOnSubscribe<String> { override fun subscribe(emitter: ObservableEmitter<String>) { emitter.onNext("somestring") emitter.onComplete() } }) .delay(10, TimeUnit.MILLISECONDS) // 延迟10ms发射,给订阅流程留时间 .subscribeOn(Schedulers.computation()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(...)
4. 检查CompositeDisposable的生命周期
确保CompositeDisposable没有在事件被主线程处理前被过早调用dispose()——比如如果test()在Activity.onCreate()中调用,而CompositeDisposable在onStop()中被dispose,那主线程消息队列里未处理的事件就会被丢弃。
额外说明
你提到“subscribe()应在订阅者就绪后才会被调用”,这个理解是完全正确的:这里的subscribe()指的是ObservableOnSubscribe中的方法,当它被执行时,订阅者已经完全绑定,可以接收事件。所以理论上每个事件都应该被触发,日志丢失的核心大概率还是日志框架的合并行为。
内容的提问来源于stack exchange,提问作者southerton81

