【问题标题】:PublishSubject `subscribeOn` behaviorPublishSubject `subscribeOn` 行为
【发布时间】:2020-08-09 20:40:01
【问题描述】:

为什么subscribe 从来不在这里打印任何东西?只是出于好奇。无论如何,这都是不好的做法:我通常会改用observeOn。但是,我无法弄清楚为什么永远无法到达 subscribe...

val subject: PublishSubject<Int> = PublishSubject.create()
val countDownLatch = CountDownLatch(1)

subject
    .map { it + 1 }
    .subscribeOn(Schedulers.computation())
    .subscribe {
        println(Thread.currentThread().name)
        countDownLatch.countDown()
    }

subject.onNext(1)
countDownLatch.await()

【问题讨论】:

  • 您确定要减少println() 语句之前的计数器吗?您的主线程可能会在 println() 语句执行之前终止。
  • 没错,那是错误的,但这不是问题所在。尝试运行该代码,你会发现它永远不会到达subscribe

标签: kotlin rx-java3


【解决方案1】:

为什么会这样

订阅的过程中,观察者通过Subscribe 通知向可观察者发出信号,表明它已准备好接收项目。详情请见Observable contract

此外,Subject 文档指出:

请注意,PublishSubject 可能会在创建后立即开始发射项目(除非您已采取措施防止这种情况发生),并且 因此在 @987654328 之间可能会丢失一个或多个项目@ 被创建并且观察者订阅它

当您尝试通过.subscribeOn(Schedulers.computation()) 订阅新线程后立即调用subject.onNext(_) 时,可观察对象(即subject)可能仍在等待来自Subscribe 的通知观察者。例如:

subject
    .subscribeOn(Schedulers.computation())
    .subscribe { println("received item") }

// this usually prints nothing!
subject.onNext(1)

但是,如果您在发出第一个项目之前添加一点时间延迟,则可观察对象更有可能在您调用 subject.onNext(_) 之前从观察者收到 Subscribe 通知。例如:

subject
    .subscribeOn(Schedulers.computation())
    .subscribe { println("received item") }

// wait for subscription to be established properly
Thread.sleep(1000)

// this usually prints "received item"
subject.onNext(1)

怎么办?

如果您希望您的所有订阅都接收可观察对象发出的所有项目,您可以执行以下操作之一:

  • 在调用subject.onNext(_)之前阻塞主线程等待所有观察者订阅。
  • 创建一个新的 observable,它会等待所有 observable 都被订阅,然后再在其内部调用 subject.onNext(_)

这些也可能有用:

  • ReplaySubject:这使您可以存储所有以前项目的历史记录,并在每次订阅时重新发送它们。缺点:您需要在内存中存储任意数量的项目。
  • ConnectableObservable:这确保了 observable 仅在调用 .connect() 之后才发出项目。特别是,.autoConnect(n) 运算符确保 observable 仅在 n 观察者成功订阅后发出。

示例:在订阅之前阻塞主线程

val subject: PublishSubject<Int> = PublishSubject.create()
val countDownLatch = CountDownLatch(1)
val isSubscribedLatch = CountDownLatch(1)

subject
    .subscribeOn(Schedulers.computation())
    .doOnSubscribe { isSubscribedLatch.countDown() }
    .map { it + 1 }
    .subscribe {
        countDownLatch.countDown()
        println(Thread.currentThread().name)
    }

isSubscribedLatch.await()
subject.onNext(1)
countDownLatch.await()

【讨论】:

    猜你喜欢
    • 2017-11-21
    • 2020-11-29
    • 2023-04-03
    • 1970-01-01
    • 2017-01-11
    • 1970-01-01
    • 1970-01-01
    • 2019-10-18
    • 1970-01-01
    相关资源
    最近更新 更多