【问题标题】:What is the Akka streams .via() equivalent in Project Reactor?Project Reactor 中的 Akka 流 .via() 等价物是什么?
【发布时间】:2021-03-31 17:46:54
【问题描述】:

我有点纠结于 Project Reactor 文档。我在 Akka Streams 方面有一些经验,但现在我正在做一个使用 Project Reactor 的项目。

我需要一个能够接受序列以传递消息的 Reactor 运算符。它的行为需要类似于 Akka Streams 中的 .via() 运算符。

例如,假设我们有一个序列:A -> B -> C,我需要在 B 步骤之后注入序列 X1 -> X2 -> X3。所以最终的顺序是 A -> B -> X1 -> X2 -> X3 -> C.

Reactor 中是否存在类似的东西?

【问题讨论】:

    标签: java scala project-reactor akka-stream


    【解决方案1】:

    Project Reactor 中接近via 的是transform 方法。

    所以在 Akka 中说你有这个图表:

    Source.single(10)
          .map(_ * -1) //some mapping
          .runWith(Sink.ignore)
    

    然后你就有了这个流程:

    val flow = Flow[Int].map(_ * 2)
    

    您可以像这样将该流程插入到您的图表中:

    Source.single(10)
          .map(_ * -1)
          .via(flow)
          .runWith(Sink.ignore)
    

    Project Reactor 中的等价物是这样的:

    有一个图表:

     Flux.just(10)
         .map(x -> x * -1)
         .subscribe();
    

    以及将Flux<Integer> 转换为Publisher<Integer> 的方法:

    public static class Transformers
    {
      public static Publisher<Integer> flow(Flux<Integer> f)
      {
        return f.map(x -> x * 2);
      }
    }
    

    您可以像这样将该方法插入到您的图表中:

    Flux.just(10)
        .map(x -> x * -1)
        .transform(Transformers::flow)
        .subscribe();
    

    我写了一篇关于这两个 API 之间的差异和其他差异的文章,也许您会发现它很有用。这篇文章来自 2019 年,API 不断发展。例如,我在Flux 的上下文中提到了compose 方法,自从我写这篇文章以来,它已重命名为transformDeferred,我不确定自从我撰写这篇文章以来还有什么漂移,所以请注意:Akka Streams vs Project Reactor API

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-08-29
      • 1970-01-01
      • 1970-01-01
      • 2014-05-08
      • 2014-06-12
      • 1970-01-01
      • 2022-11-28
      • 2021-06-19
      相关资源
      最近更新 更多