【问题标题】:RxJava2 PublishSubject subscriber fails to receive items when called from multiple threads using SingleScheduler使用 SingleScheduler 从多个线程调用时,RxJava2 PublishSubject 订阅者无法接收项目
【发布时间】:2017-03-24 18:13:00
【问题描述】:

我有以下单元测试,我尝试从不同的线程发送 10 个Strings,并测试我从单个线程接收到这些Strings。我的问题是这个测试失败了。有时它会成功,但有时我只收到89 项目,然后测试挂起直到闩锁超时。我是否以错误的方式使用SingleScheduler?我错过了什么吗?

val consumerCallerThreadNames = mutableSetOf<String>()
val messageCount = AtomicInteger(0)

val latch = CountDownLatch(MESSAGE_COUNT)

@Test
fun someTest() {
    val msg = "foo"

    val subject = PublishSubject.create<String>()
    subject
            .observeOn(SingleScheduler())
            .subscribe({ message ->
                consumerCallerThreadNames.add(Thread.currentThread().name)
                messageCount.incrementAndGet()
                latch.countDown()
            }, Throwable::printStackTrace)

    1.rangeTo(MESSAGE_COUNT).forEach {
        Thread({
            try {
                subject.onNext(msg)
            } catch (t: Throwable) {
                t.printStackTrace()
            }
        }).start()
    }
    latch.await(10, SECONDS)

    assertThat(consumerCallerThreadNames).hasSize(1)
    assertThat(messageCount.get()).isEqualTo(MESSAGE_COUNT)
}

companion object {
    val MESSAGE_COUNT = 10
}

如果我将其重写为使用单线程 ExecutorService,则测试不再出现问题,因此问题出在 Rx 或我对 Rx 缺乏了解。

【问题讨论】:

    标签: java multithreading kotlin rx-java2


    【解决方案1】:

    RxJava 要求不能同时调用on*。这意味着您的代码不是线程安全的。

    由于只有主题本身以并发方式使用,它应该可以通过使用 Subject&lt;T&gt;.toSerialized() 方法序列化(本质上是 Java 的“同步”)主题本身来修复。

    val subject = PublishSubject.create&lt;String&gt;() 变为 val subject = PublishSubject.create&lt;String&gt;().toSerialized()

    【讨论】:

    • same time 是什么意思?我不能快速连续拨打onNext 吗?这意味着 Rx 本身不是线程安全的。例如,在现实生活中可能会发生两个用户同时单击一个按钮,我通过这种方式得到两个 onNexts。
    • Rx 规定,任何调用Observer 接口的人都必须确保没有并发调用。为了支持不可能做到这一点的情况,toSerialized() 运营商会以一些开销强制执行它。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多