【问题标题】:How can I pipe the output of an Akka Streams Merge to another Flow?如何将 Akka Streams Merge 的输出通过管道传输到另一个 Flow?
【发布时间】:2014-12-01 16:43:02
【问题描述】:

我正在使用 Akka Streams,并且已经了解了大部分基础知识,但我不清楚如何获取 Merge 的结果并进行进一步的操作(映射、过滤、折叠等)在上面。

我想修改以下代码,以便我可以进一步操作数据,而不是将合并通过管道传输到接收器。

implicit val materializer = FlowMaterializer()

val items_a = Source(List(10,20,30,40,50))
val items_b = Source(List(60,70,80,90,100))
val sink = ForeachSink(println)

val materialized = FlowGraph { implicit builder =>
  import FlowGraphImplicits._
  val merge = Merge[Int]("m1")
  items_a ~> merge
  items_b ~> merge ~> sink
}.run()

我想我的主要问题是我不知道如何制作一个没有源的流组件,而且我不知道如何在不使用特殊的 Merge 对象和~> 语法。

编辑:这个问题和答案适用于 Akka Streams 0.11 并与之合作

【问题讨论】:

    标签: scala stream akka


    【解决方案1】:

    如果您不关心 Merge 的语义,其中元素随机向下游移动,那么您可以在 Source 上尝试 concat,而不是像这样:

    items_a.concat(items_b).map(_ * 2).map(_.toString).foreach(println)
    

    这里的区别是来自a 的所有项目将在b 的任何元素之前先向下游流动。如果你真的需要Merge 的行为,那么你可以考虑如下(请记住,你最终需要一个接收器,但你当然可以在合并后进行额外的转换):

    val items_a = Source(List(10,20,30,40,50))
    val items_b = Source(List(60,70,80,90,100))
    
    val sink = ForeachSink[Double](println)
    val transform = Flow[Int].map(_ * 2).map(_.toDouble).to(sink)
    
    
    val materialized = FlowGraph { implicit builder =>
      import FlowGraphImplicits._
      val merge = Merge[Int]("m1")
      items_a ~> merge
      items_b ~> merge ~> transform
    }.run
    

    在此示例中,您可以看到我使用来自 Flow 伴侣的帮助程序来创建一个 Flow,而无需特定输入 Source。然后我可以从那里将其附加到合并点以进行额外的处理。

    【讨论】:

    • 谢谢,这正是我想要的。
    【解决方案2】:

    使用Source.combine:

    val items_a :: items_b :: items_c = List(
             Source(List(10,20,30,40,50)), 
             Source(List(60,70,80,90,100), 
             Source(List(110,120,130,140,1500))
    
    Source.combine(items_a, items_b, items_c : _*)(Merge(_))
             .map(_+1)
             .runForeach(println)
    

    【讨论】:

      【解决方案3】:

      或者,如果您需要保留输入源的顺序(例如 items_a 必须在 items_b 之前,并且 items_b 必须在 items_c 之前),您可以使用 Concat,而不是 Merge。

      val items_a :: items_b :: items_c = List(
           Source(List(10,20,30,40,50)), 
           Source(List(60,70,80,90,100), 
           Source(List(110,120,130,140,1500))
      Source.combine(items_a, items_b, items_c : _*)(Concat(_))
      

      【讨论】:

        猜你喜欢
        • 2018-11-17
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-09-30
        • 2016-02-13
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多