RxJava中带debounce的Observable如何合理关闭Subscriber?
嘿,作为RxJava新手碰到这个问题太正常了,我来帮你理清楚怎么解决~
首先先指出你代码里的一个小问题:你是先调用subscribe(mySubscriber)再定义mySubscriber的,这会导致编译错误,得先把订阅者定义好再去订阅流哈。
接下来聊核心问题:为什么onCompleted没触发?因为你的原Observable是持续从控制器接收数据的,它本身不会主动发送完成信号;debounce只是过滤事件,也不会帮你触发完成。而你需要的是:在最后一次debounce输出事件后,启动一个超时计时器,如果指定时间(比如20秒)内没有新的debounce输出,就触发onCompleted关闭订阅;但如果期间有新的debounce输出,就重置这个计时器。
解决方案:结合debounce + timeout操作符
我们可以用timeout的重载版本,指定超时时发送一个空Observable(Observable.empty()),这样超时后就会触发完成信号,而不是错误信号。具体实现如下:
import rx.Observable; import rx.Subscriber; import java.util.concurrent.TimeUnit; // 先修正订阅者定义顺序,再添加超时逻辑 final Subscriber<Long> mySubscriber = new Subscriber<Long>() { @Override public void onCompleted() { System.out.println("killing the subscriber"); } @Override public void onError(final Throwable throwable) { // 建议这里加错误处理,避免意外终止订阅 System.err.println("订阅出错:" + throwable.getMessage()); } @Override public void onNext(final Long number) { // 处理你的业务逻辑 System.out.println("处理事件:" + number); } }; observable.asObservable() .debounce(10, TimeUnit.SECONDS) // 原逻辑:10秒内无新事件则输出最后一个 .timeout(20, TimeUnit.SECONDS, Observable.empty()) // 关键逻辑:最后一次debounce输出后,20秒无新事件则触发完成 .subscribe(mySubscriber);
逻辑验证(对应你的例子)
咱们用你举的场景走一遍流程:
- 第5秒收到第一个事件 →
debounce开始计时到第15秒 → 第15秒输出该事件 →timeout计时器启动,到第35秒超时; - 第14秒收到第二个事件 →
debounce计时器重置,到第24秒输出该事件 →timeout计时器同步重置,到第44秒超时; - 如果之后没有新事件,第44秒时
timeout触发,发送onCompleted信号,订阅者正常关闭,完全不会影响第24秒的事件处理。
如果在超时前(比如第30秒)又收到新事件,debounce会再次重置计时,输出新事件后timeout也会跟着重置计时器,完美符合你的需求。
额外注意点
- 如果原Observable后续可能重新发送事件,而你需要重新订阅的话,可以考虑结合
retry操作符,或者在onCompleted里重新创建订阅; - 记得在合适的时机调用
mySubscriber.unsubscribe(),避免内存泄漏。
内容的提问来源于stack exchange,提问作者jpganz18
相关产品推荐
相关产品推荐

