如何确保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) } }
关键说明:
startAndForget:该方法会让serviceB的Task在后台调度执行,无需等待完成,也不会处理其返回值或错误。如果需要处理serviceB的错误,可以用start监听结果,但依然不阻塞主流程。- 调度器隔离:
serviceB依旧使用bCtx调度器执行,和serviceA的aCtx互不干扰,保证资源隔离。 - 并行性不受影响:
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
相关产品推荐
相关产品推荐

