【问题标题】:Reactor on-demand flux or a sink反应堆按需通量或水槽
【发布时间】:2021-08-10 09:29:49
【问题描述】:

考虑一个 HTTP 控制器端点,它接受请求、验证然后返回 ack,但同时在后台做一些“繁重的工作”。

2 种方法使用 Reactor(我感兴趣)可以实现:

第一种方法

@PostMapping(..)
fun acceptRequest(request: Request): Response {
  if(isValid(request)) {
    Mono.just(request)
      .flatMap(service::doHeavyWork)
      .subscribe(...)
    return Response(202)
  } else {
    return Response(400)
  }
}

第二种方法

class Controller {
  private val service = ...
  private val sink = Sinks.many().unicast().onBackpressureBuffer<Request>()
  private val stream = sink.asFlux().flatMap(service::doHeavyWork).subscribe(..)


  fun acceptRequest(request: Request): Response {
    if(isValid(request)) {
      sink.tryEmitNext(request) // for simplicity not handling errors
      return Response(202)
    } else {
      return Response(400)
    }
  }
}

哪种方法更好/更差,为什么?

我问的原因是,在 Akka 中,按需构建流并不是真正有效的,因为流需要每次都实现,所以最好采用“接收器方法”。我想知道这是否也适用于反应堆,或者使用这些方法是否还有其他优点/缺点。

【问题讨论】:

    标签: spring-boot reactive-programming project-reactor


    【解决方案1】:

    我对 Akka 不太熟悉,但是使用 Reactor 构建响应式链肯定不会带来巨大的开销——这是处理请求的“正常”方式。因此,我认为不需要像您的第二种方法那样使用单独的接收器 - 这似乎只是增加了复杂性而收效甚微。因此,我会说第一种方法更好。

    话虽如此,一般来说,不建议像在这两个示例中那样订阅自己 - 但这种“一劳永逸”的工作是它可能有意义的少数情况之一。我在这里提出的其他几个潜在警告可能值得考虑:

    • 你称这项工作为“繁重”,我不确定这是否意味着它的 CPU 很重,或者只是 IO 很重,或者需要很长时间。如果只是因为触发了一堆请求而需要很长时间,那没什么大不了的。但是,如果它的 CPU 很重,那么如果您不小心,可能会导致问题 - 您可能不想在事件循环线程上执行 CPU 繁重的任务。在这种情况下,我可能会创建一个由您自己的执行程序服务支持的单独调度程序,然后使用subscribeOn() 来卸载那些 CPU 繁重的任务。
    • 请记住,在这种情况下,“一劳永逸”模式实际上是“忘记”——您完全无法知道您卸载的繁重任务是否有效,因为您基本上已经将这些信息扔掉了自行订阅。根据您的用例,这可能没问题,但如果任务很关键,或者如果任务失败您需要某种反馈,那么值得考虑的是,这可能不是最好的方法。

    【讨论】:

    • 主要是 IO 并且可能需要一些时间……而不是 CPU 密集型。是的,我知道即发即弃有一些缺点。
    • IO 重绝对没问题 - 这非常适合反应式,因为您保持打开的连接通常会在阻塞模型中占用线程,但在反应式模型中基本上是免费的。
    猜你喜欢
    • 2017-10-17
    • 2023-03-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-02-04
    • 2016-07-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多