【问题标题】:Consume a source with two sinks and get the result of one sink用两个 sink 消费一个 source 并得到一个 sink 的结果
【发布时间】:2019-07-09 15:06:21
【问题描述】:

我想使用带有两个不同接收器的Source

简化示例:

val source = Source(1 to 20)

val addSink = Sink.fold[Int, Int](0)(_ + _)
val subtractSink = Sink.fold[Int, Int](0)(_ - _)

val graph = GraphDSL.create() { implicit builder =>
  import GraphDSL.Implicits._

  val bcast = builder.add(Broadcast[Int](2))

  source ~> bcast.in

  bcast.out(0) ~> addSink
  bcast.out(1) ~> subtrackSink

  ClosedShape
}

RunnableGraph.fromGraph(graph).run()

val result: Future[Int] = ???

我需要能够检索addSink 的结果。 RunnableGraph.fromGraph(graph).run() 给我NotUsed, 但我想得到一个Int(第一次折叠的结果Sink)。有可能吗?

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    将两个接收器传递给图形构建器的 create 方法,这使您可以访问它们各自的具体化值:

    val graph = GraphDSL.create(addSink, subtractSink)((_, _)) { implicit builder =>
      (aSink, sSink) =>
      import GraphDSL.Implicits._
    
      val bcast = builder.add(Broadcast[Int](2))
    
      source ~> bcast.in
      bcast.out(0) ~> aSink
      bcast.out(1) ~> sSink
      ClosedShape
    }
    
    val (addResult, subtractResult): (Future[Int], Future[Int]) =
      RunnableGraph.fromGraph(graph).run() 
    

    或者,您可以放弃图形 DSL 并使用alsoToMat

    val result: Future[Int] =
      Source(1 to 20)
        .alsoToMat(addSink)(Keep.right)
        .toMat(subtractSink)(Keep.left)
        .run()
    

    上面给出了addSink 的具体化值。如果要获取addSinksubtractSink 的具体化值,请使用Keep.both

    val (addResult, subtractResult): (Future[Int], Future[Int]) =
      Source(1 to 20)
        .alsoToMat(addSink)(Keep.right)
        .toMat(subtractSink)(Keep.both) // <--
        .run()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-08-31
      • 2020-03-25
      • 1970-01-01
      • 1970-01-01
      • 2021-11-13
      • 1970-01-01
      • 2017-01-08
      • 2019-03-01
      相关资源
      最近更新 更多