手动终止takeUntil内的Observable.timer不生效问题如何解决
问题核心原因分析
- 你单独调用
Observable.timer(5, TimeUnit.SECONDS).also { disposable = it.subscribe() }拿到的Disposable,和takeUntil运算符内部订阅该timer生成的订阅实例完全独立。你dispose自己持有这个实例,不会影响takeUntil内部持有的timer订阅,timer仍会正常走完5秒倒计时才会触发终止逻辑。 - 就算你能dispose
takeUntil内部的timer订阅,反而会导致timer永远不会发射终止事件,takeUntil永远等不到终止信号,事件流会一直运行,完全违背需求。takeUntil的终止逻辑依赖收到传入Observable的发射事件/终止事件,而非传入的Observable被销毁。
正确实现方案
你的需求是同时支持5秒自动终止和手动触发终止两个结束条件,应该用PublishSubject作为手动终止的信号发射器,将两个终止条件合并后传入takeUntil即可。
代码示例
首先声明成员变量:
// 手动终止信号发射器 private val stopSignal = PublishSubject.create<Unit>() // 持有整个流的订阅实例,用于页面销毁时整体释放资源 private var streamDisposable: Disposable? = null
核心流实现:
streamDisposable = observableThatHasToBeAliveAllTime .switchMap { observableThatEmitsItemOver5SecsWhenUpperObsEmits() .takeUntil( // 两个终止条件任意触发一个,就终止当前事件流 Observable.ambArray( Observable.timer(5, TimeUnit.SECONDS), stopSignal.take(1) ) ) .switchMap { /* some work */ } } .subscribe { /* handle result */ }
手动终止时调用以下代码即可立即终止当前运行的事件流:
stopSignal.onNext(Unit)
注意事项
- 若
observableThatHasToBeAliveAllTime为全局存活的热流,页面销毁时需调用streamDisposable?.dispose()释放资源,避免内存泄漏。 - 该实现支持多次触发5秒监听窗口和多次手动终止,无需重建
stopSignal实例。
内容的提问来源于stack exchange,提问作者lishee loo
相关产品推荐
相关产品推荐

