RxJava:合并冷热Observable实现相互等待的问题求助
问题分析与解决方案
这个问题我之前也碰到过,核心原因其实很简单:concatWith的工作机制是必须等前一个Observable完全结束后,才会去订阅后一个Observable。而RxView.clicks()是个典型的冷Observable——在它被订阅之前发生的点击事件,根本不会被记录下来,所以你在initLoading延迟期间点的按钮,自然就“丢”了。
解决方案:将点击事件转为可缓存的热Observable
我们需要把按钮点击的事件变成能提前收集、缓存的热Observable,让它在initLoading还在执行的时候就开始监听点击,把事件存起来,等initLoading结束后再把缓存的事件和后续新事件一起发射给订阅者。
方案1:使用replay() + connect()
这是最简洁的写法,直接利用RxJava的操作符实现:
val WAIT_TIME = 10L // 假设你的延迟时长为10秒 val initLoading = Observable.fromCallable { println("${System.currentTimeMillis()}") } .subscribeOn(Schedulers.computation()) .delay(WAIT_TIME, TimeUnit.SECONDS) .map { "loading ${System.currentTimeMillis()}" } .observeOn(AndroidSchedulers.mainThread()) // 将点击事件转为可连接Observable,replay()会缓存所有订阅前的事件 val click = RxView.clicks(button) .map { "click ${System.currentTimeMillis()}" } .replay() // 提前连接,让Observable立刻开始监听并缓存点击事件 val clickConnection = click.connect() initLoading.concatWith(click) .subscribeBy( onNext = { println("result $it") }, onError = { throw it } ) // 页面销毁时释放资源,避免内存泄漏 override fun onDestroy() { super.onDestroy() clickConnection.dispose() }
原理说明
replay()创建的ConnectableObservable,在调用connect()后会立即开始订阅上游的点击事件,并缓存所有发射的事件。- 当
initLoading执行完成后,concatWith会订阅这个clickObservable,此时它会先发射缓存的点击事件(你在延迟期间的点击),之后再继续发射新的点击事件。
如果只需要保留最近一次点击(比如用户连续点击多次,仅需最后一次),可以改为replay(1),这样只会缓存最近1个事件,更节省内存。
方案2:使用ReplaySubject/BehaviorSubject
如果你更习惯用Subject来管理事件流,也可以直接使用ReplaySubject(缓存所有事件)或BehaviorSubject(缓存最近一次事件):
val clickSubject = ReplaySubject.create<String>() // 让点击事件流入Subject RxView.clicks(button) .map { "click ${System.currentTimeMillis()}" } .subscribe(clickSubject) initLoading.concatWith(clickSubject) .subscribeBy( onNext = { println("result $it") }, onError = { throw it } ) // 页面销毁时记得结束Subject,避免内存泄漏 override fun onDestroy() { super.onDestroy() clickSubject.onComplete() }
这个思路和replay()完全一致,因为replay()内部就是基于ReplaySubject实现的,只是操作符的写法更符合RxJava的链式风格。
内容的提问来源于stack exchange,提问作者Rezo Shalikashvili
相关产品推荐
相关产品推荐

