You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

为什么这样有效

  1. defer的作用:延迟Observable的创建逻辑,直到有订阅发生。当所有订阅取消、refCount降到0后,下次新订阅到来时,defer会重新执行内部代码,创建一个全新的Observable实例,从而触发日志打印和数据源重启。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 13:33:40