【问题标题】:Can I notify a BehaviorProcessor from inside a RxJava stream?我可以从 RxJava 流中通知 BehaviorProcessor 吗?
【发布时间】:2019-01-16 16:05:18
【问题描述】:

我想获得您对以下代码的反馈。 我想知道拨打currentSession.onNext(result.session) 是否安全 从SessionManager.signIn 流内部。

由于多线程和同步问题,我的第一直觉是说NO,这意味着,根据这段代码,我可以从不同的线程调用currentSession.onNext(result.session)

这里是代码,请告诉我你的想法!谢谢

SessionManager 是一个单例

@Singleton
class SessionManager @Inject constructor(
    private val sessionService: SessionService,
){

    val currentSession = BehaviorProcessor.create<Session>()

    fun signIn(login: String, password: String): Single<Boolean> =
        sessionService.signIn(login, password)
            .doOnNext(result -> 
                if (session is Success) {
                   currentSession.onNext(result.session)
                }
            ).map { result ->
                when (result) {
                    is Success -> true
                    else -> false
                }
            }
            .subscribeOn(Schedulers.io())
}

HomeView 是订阅 SessionManager 的登录流的随机视图

class HomeView(val context: Context) : View(context) {

        @Inject
        lateinit var sessionManager: SessionManager

        private val disposables = CompositeDisposable()

        override fun onAttachedToWindow() {
            super.onAttachedToWindow()

            disposables.add(sessionManager.signIn("username", "password")
                .distinctUntilChanged()
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe { result ->
                    textView.text = if (result) "Success" else "Fail"
                })
        }

        override fun onDetachedFromWindow() {
            super.onDetachedFromWindow()
            disposables.clear()
        }
    }

SessionManager观察currentSession的随机视图

class RandomView(val context: Context) : View(context) {

        @Inject
        lateinit var sessionManager: SessionManager

        private val disposables = CompositeDisposable()

        override fun onAttachedToWindow() {
            super.onAttachedToWindow()

            disposables.add(sessionManager.currentSession
                .distinctUntilChanged()
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe { session -> userTextView.text = session.userName })
        }

        override fun onDetachedFromWindow() {
            super.onDetachedFromWindow()
            disposables.clear()
        }
    }

【问题讨论】:

    标签: multithreading kotlin rx-java2 observer-pattern


    【解决方案1】:

    documentation of BehaviorProcessor 说:

    调用 onNext(Object)、offer(Object)、onError(Throwable) 和 onComplete() 需要被序列化(从同一个线程调用或通过外部序列化方式从不同线程非重叠调用)。 所有 FlowableProcessor 都可用的 FlowableProcessor.toSerialized() 方法提供了这样的序列化并防止重入(即,当使用此处理器的下游订阅者也希望在此处理器上递归调用 onNext(Object) 时)。

    所以如果你这样定义它:

    val currentSession = BehaviorProcessor.create<Session>().toSerialized()
    

    那么你可以安全地从任何线程调用onNext,它不会引起任何同步问题。

    注意事项:

    我同意处理器的更新应该在doOnNext而不是map

    我认为最好使用Completable 而不是Single&lt;Boolean&gt;,并使用Rx 错误来指示阻止登录的原因。您还应该在subscribe 方法中定义错误处理程序。

    【讨论】:

    • 有道理 - 谢谢!之后我发现了一篇关于它的好文章:proandroiddev.com/… 他们准确地提到了你在回复中所说的话。也感谢您的旁注,我自愿保持流尽可能简单以避免噪音,但您的建议很有帮助!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-04-16
    • 2012-08-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多