【问题标题】:Multiple sinks in the same stream同一流中的多个接收器
【发布时间】:2017-12-19 22:17:11
【问题描述】:

我有一个这样的流和两个接收器,但一次只使用一个:

Source.fromElements(1, 2, 3)
.via(flow)
.runWith(sink1)

Source.fromElements(1, 2, 3)
.via(flow)
.runWith(sink2)

我们使用哪个接收器是可配置的,但如果我同时使用两个接收器会怎样。 我该怎么做?

我考虑过 Sink.combine,但它还需要一个合并策略,我不想以任何方式组合这些接收器的结果。我并不真正关心它们,所以我只想通过 HTTP 将相同的数据发送到某个端点,同时将它们发送到数据库。 Sink combine 与广播非常相似,但从头开始实现广播会降低代码的可读性,现在我只有简单的源、流和接收器,没有低级图形阶段。

你知道如何做到这一点的正确方法吗(有背压和其他我只使用一个水槽的东西)?

【问题讨论】:

    标签: scala akka akka-stream reactive-streams


    【解决方案1】:

    你可以使用alsoTo(见API docs):

    Flow[Int].alsoTo(Sink.foreach(println(_))).to(Sink.ignore)
    

    【讨论】:

    • 如何通过在第二个接收器之前添加一个简单的 .async 来并行运行这些接收器?我想并行运行它们,但仍然有背压,换句话说,我希望我的流运行速度与在最慢接收器中花费的时间一样快,而不是在所有接收器中花费的时间总和(因为它们是同步运行的)。
    • 值得一提的是,在 to() 中声明了 sink 之后,alsoTo() 中的 sink 将在最后执行。
    【解决方案2】:

    以最简单的形式使用GraphDSL 进行广播不应降低可读性——事实上,人们甚至可能会争辩说~> 子句在某种程度上有助于可视化流结构:

    val graph = RunnableGraph.fromGraph(GraphDSL.create() { implicit builder =>
      import GraphDSL.Implicits._
      val bcast = builder.add(Broadcast[Int](2))
    
      Source.fromElements(1, 2, 3) ~> flow ~> bcast.in
      bcast.out(0) ~> sink1
      bcast.out(1) ~> sink2
    
      ClosedShape
    })
    graph.run()
    

    【讨论】:

    • 这些接收器是否并行运行?
    • 默认情况下,Akka Streams 顺序执行图形处理阶段,但如果需要,您可以使用方法async 并行执行它们。欲了解更多详情,请点击Akka Stream doc 讨论该主题。
    猜你喜欢
    • 1970-01-01
    • 2012-03-04
    • 1970-01-01
    • 2018-06-10
    • 1970-01-01
    • 2018-10-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多