• 首页 首页 icon
  • 工具库 工具库 icon
    • IP查询 IP查询 icon
  • 内容库 内容库 icon
    • 快讯库 快讯库 icon
    • 精品库 精品库 icon
    • 问答库 问答库 icon
  • 更多 更多 icon
    • 服务条款 服务条款 icon

容易被忽略的知识点RxJava操作符的线程安全,kotlin源码

武飞扬头像
m0_66265001
帮助1

val publishSubject = PublishSubject.create()

val actuallyReceived = AtomicInteger()

publishSubject.take(3).subscribe {

actuallyReceived.incrementAndGet()

}

val latch = CountDownLatch(numberOfThreads)

var threads = listOf()

(0…numberOfThreads).forEach {

threads = thread(start = false) {

publishSubject.onNext(it)

latch.countDown()

}

}

threads.forEach { it.start() }

latch.await()

check(actuallyReceived.get() == 3)

}

}

执行上面代码,由于take的结果不符合预期,总是会异常退出

学新通

看一下take的源码:

public final class ObservableTake extends AbstractObservableWithUpstream<T, T> {

final long limit;

public ObservableTake(ObservableSource source,

这篇好文章是转载于:学新通技术网

  • 版权申明: 本站部分内容来自互联网,仅供学习及演示用,请勿用于商业和其他非法用途。如果侵犯了您的权益请与我们联系,请提供相关证据及您的身份证明,我们将在收到邮件后48小时内删除。
  • 本站站名: 学新通技术网
  • 本文地址: /boutique/detail/tanhgccaba
系列文章
更多 icon
同类精品
更多 icon
继续加载