同一个流中有多个接收器

Pio*_*ski 4 scala akka reactive-streams akka-stream

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

Source.fromElements(1, 2, 3)
.via(flow)
.runWith(sink1)
Run Code Online (Sandbox Code Playgroud)

要么

Source.fromElements(1, 2, 3)
.via(flow)
.runWith(sink2)
Run Code Online (Sandbox Code Playgroud)

它是可配置的我们使用的接收器,但是如果我并行使用两个接收器怎么办呢.我怎样才能做到这一点?

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

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

Seb*_*ian 10

您可以使用alsoTo(请参阅API文档):

Flow[Int].alsoTo(Sink.foreach(println(_))).to(Sink.ignore)
Run Code Online (Sandbox Code Playgroud)

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