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

如何确保Monix Observable中serviceB延迟不影响serviceA的流式处理?

版本信息:

Scala版本:2.12.9,monix/eval:2.3.2,monix/reactive:2.3.2

背景:

需要从Monix Observable中调用serviceA和serviceB两个服务,程序代码如下:

import monix.reactive.subjects.PublishToOneSubject
import monix.eval.Task
import monix.execution.Scheduler
import monix.execution.Ack.Continue
import java.time.LocalDateTime
import java.time.format.DateTimeFormatter

val formatter: DateTimeFormatter = DateTimeFormatter.ofPattern("HH:mm:ss.SSSS");
def timestamp = formatter.format(LocalDateTime.now)
def log = s"[$timestamp] [${Thread.currentThread.getName}]"

println(s"$log Program Starts.")

private def serviceA(num: Int) = Task {
  println(s"$log A: ${num}")
  Thread.sleep(1000)
  num * 10
}

private def serviceB(num: Int) = Task {
  println(s"$log B: ${num}")
  Thread.sleep(5 * 1000)
  num * 0.5
}

private val aCtx = Scheduler.fixedPool("aCtx", 2)
private val bCtx = Scheduler.fixedPool("bCtx", 2)

val stream: PublishToOneSubject[Int] = PublishToOneSubject[Int]

stream
  .doOnSubscribe( () => println(s"$log Stream is ready!"))
  .doOnNext {
    v => println(s"$log Event: ${v}")
  }
  .mapAsync(parallelism=3) {
    v =>
      serviceA(v)
        .zip(serviceB(v).executeOn(bCtx))
        .map(_._1)
  }
  .subscribe(
    nextFn = { v =>
        println(s"$log Result: ${v}")
        Continue
    }
  )(aCtx)

(1 to 10).foreach {
  i => stream.onNext(i)
}


Thread.sleep(30 * 1000)
println(s"$log Program Ends.")

问题:

如何确保调用serviceB时的延迟不会影响流式处理流程,让流可以持续处理事件并返回serviceA的结果?

尽管需要调用serviceB,但仅关注serviceA的结果,不希望serviceB的慢调用拖慢serviceA的处理。当前使用zip会因serviceB的延迟导致结果返回延迟。

解决方案:

核心思路是让serviceB在后台异步执行,无需等待它完成,这样serviceA的结果可以立刻返回给流,不会被serviceB的延迟阻塞。

修改方式:

把原来的zip逻辑替换为:先执行serviceA,同时触发serviceB的后台执行但不等待其结果。具体代码调整如下:

.mapAsync(parallelism=3) {
  v =>
    // 执行serviceA,同时后台启动serviceB
    serviceA(v).flatMap { aResult =>
      // 启动serviceB异步执行,不等待结果
      serviceB(v).executeOn(bCtx).startAndForget
      // 直接返回serviceA的结果
      Task.pure(aResult)
    }
}

关键说明:

  1. startAndForget:该方法会让serviceB的Task在后台调度执行,无需等待完成,也不会处理其返回值或错误。如果需要处理serviceB的错误,可以用start监听结果,但依然不阻塞主流程。
  2. 调度器隔离:serviceB依旧使用bCtx调度器执行,和serviceA的aCtx互不干扰,保证资源隔离。
  3. 并行性不受影响:mapAsync的parallelism=3仍控制serviceA的并行处理数量,serviceB的并行由bCtx的线程池大小(2)决定,两者互不影响。

完整修改后的核心代码片段:

stream
  .doOnSubscribe( () => println(s"$log Stream is ready!"))
  .doOnNext {
    v => println(s"$log Event: ${v}")
  }
  .mapAsync(parallelism=3) {
    v =>
      serviceA(v).flatMap { aRes =>
        // 后台执行serviceB,不阻塞主流程
        serviceB(v).executeOn(bCtx).startAndForget
        Task.pure(aRes)
      }
  }
  .subscribe(
    nextFn = { v =>
        println(s"$log Result: ${v}")
        Continue
    }
  )(aCtx)

修改后,serviceA的结果会在1秒后立刻返回给流,serviceB则在后台独立执行5秒,完全不会拖慢主流程的处理速度。

内容的提问来源于stack exchange,提问作者Nazil Khan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:59:52