【发布时间】:2019-04-22 16:40:30
【问题描述】:
当我创建 5 个可观察对象并使用单独的订阅者订阅每个对象时,直觉上我认为每个订阅者都会获得其可观察对象的相应数据,通过 onNext() 调用发出:
val compositeSubscription = CompositeDisposable()
fun test() {
for (i in 0..5) {
compositeSubscription.add (Observable.create<String>(object : ObservableOnSubscribe<String> {
override fun subscribe(emitter: ObservableEmitter<String>) {
emitter.onNext("somestring")
emitter.onComplete()
}
}).subscribeOn(Schedulers.computation())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({
Logger.i("testIt onNext")
}, {
Logger.i("testIt onError")
}))
}
}
但是,我在日志中看到的是一两个“testIt onNext”。
现在,当我在订阅者的 onNext() 中添加延迟时,所有 6 个订阅者 onNext() 都会被调用。
当一些订阅者的速度不够快而无法赶上他们的数据时,这似乎是一种不合理的情况。我想知道这是怎么发生的,因为 subscribe() 应该在订阅者启动并运行后调用。
如果有任何提示,我们将不胜感激。
【问题讨论】:
-
你可能不想订阅主线程
-
是的,这是真的,改为 Schedulers.computation(),现在在主线程上观察,但仍然是相同的日志。
-
从这段代码来看,每个订阅者都应该打印“testIt onNext”。你确定它没有被打印?也许 Android Studio 正在折叠相同的行?您是否尝试过为每个订阅者打印不同的内容?
-
@gpunto 你是绝对正确的,logcat 折叠相同的行(具有相同的时间戳),而且它发生得很快,这就是我的情况。如果您将其写为答案,我会接受,谢谢。