基于Monix、Scala、Cask实现Warm Observable及异常排查
问题分析与解决方案
为什么现有代码无法触发新订阅重启
你用的share.publish.refCount(Monix中share等价于publish.refCount(1))逻辑本身没问题,但核心问题是原Observable是一次性执行的:当最后一个订阅取消、refCount降到0时,原Observable已经执行完成并进入终止状态,后续新订阅只会收到完成信号,不会重新执行数据源初始化逻辑(也就不会打印SS Warm Observable日志)。
正确的Warm Observable实现方式
要满足你提出的4个需求,关键是让Observable在订阅数从0恢复到1时,能重新初始化数据源。需要结合Observable.defer和share来实现:
代码示例
假设你原来的Observable创建逻辑是这样的:
// 错误写法:Observable创建后一次性执行,无法重启 val badWarmObs = Observable.create[String] { subscriber => println("SS Warm Observable") // 模拟产生数据的逻辑 subscriber.onNext("data") Cancelable(() => println("数据源停止")) }.share
改成正确的实现:
// 正确写法:用defer包裹,每次重启都重新创建Observable val correctWarmObs = Observable.defer { Observable.create[String] { subscriber => println("SS Warm Observable") // 模拟持续产生数据的逻辑(比如定时推送) val task = Observable.interval(1.second).map(_.toString).subscribe(subscriber) Cancelable { () => task.cancel() println("数据源停止") } } }.share
为什么这样有效
defer的作用:延迟Observable的创建逻辑,直到有订阅发生。当所有订阅取消、refCount降到0后,下次新订阅到来时,defer会重新执行内部代码,创建一个全新的Observable实例,从而触发日志打印和数据源重启。share的作用:- 第一个订阅到来时自动连接数据源(满足“至少一个订阅才产生数据”)
- 所有订阅共享同一个数据源实例(满足“共享结果,不重复执行代码”)
- 最后一个订阅取消时自动断开数据源(满足“所有订阅关闭时停止产生数据”)
额外注意点
确保WebSocket通道关闭时正确取消订阅,否则refCount不会降到0,数据源不会停止,新订阅也不会触发重启。比如在Cask中的处理逻辑:
cask.websocket("/ws") { req => var subscription: Cancelable = Cancelable.empty val observer = new Observer[String] { override def onNext(value: String): Unit = req.send(value) override def onError(ex: Throwable): Unit = req.close() override def onComplete(): Unit = req.close() } subscription = correctWarmObs.subscribe(observer) cask.WsHandler { case cask.WsClose(_, _) => subscription.cancel() // 必须调用取消,让refCount减1 } }
内容的提问来源于stack exchange,提问作者it2901
相关产品推荐
相关产品推荐

