【问题标题】:Multiplex Akka sources/flows based on condition根据条件复用 Akka 源/流
【发布时间】:2021-02-11 10:27:48
【问题描述】:

有没有办法根据一些外部条件复用两个或多个 Akka 源或流?它可能看起来像这样:

def cond: Boolean = ???

val src1 = Source.fromIterator(i1)
val src2 = Source.fromIterator(i2)
val src3 = Source.mux(src1, src2, cond)

取决于cond 结果src3 应该包含来自src1 的项目或来自src2 的项目,不能同时包含两者。

我发现似乎是相反的操作divertTo。同时,似乎没有一个扇入操作支持条件合并。

【问题讨论】:

  • 什么意思?您能否添加示例输入和所需的输出?
  • 不确定还需要什么示例。假设有两个独立的事件源,我想根据通常应该在事件本身外部的动态条件,将它们组合成只有一个被发送到下游。
  • 假设src1 = Source(1, 2, 3)src2 = Source(2, 3, 4)cond = _ % 2 == 0 所以src3 = Source(2, 2, 4) 你期待这样的事情吗?
  • 没有。我的意思是,不一定。条件不必依赖于源中的项目。假设,我们有两个独立的无限流,A 和 B。假设每天下午 5 点到 7 点之间,我只想接收来自流 A 的消息,而剩下的时间 - 来自流 B。事实上现在我考虑了一下,也许一种方法是用一些 id 压缩每个源,然后合并它们,然后通过这个 id 过滤结果源。我不确定这是否是正确的方法。
  • 所以请像您在上一条评论中所做的那样详细说明,并创建minimal reproducible example,以便我们更好地了解如何回答您的问题。

标签: scala akka akka-stream


【解决方案1】:

我会建议以下内容:

def mux[T](a: Source[T, Any], b: Source[T, Any])(cond: Int => Boolean): Source[T, Any] = {
  a.map((1, _)).merge(b.map((2, _)))
    .filter(t => cond(t._1))
    .map(_._2)
}

简单地说,它为每个源发出的元素附加一个标识符(这里是 1、2,但可以是任何东西),然后使用提供的 cond 函数过滤以仅保留来自当前选定源的元素,然后映射回元素。

我认为zip 不是一个好主意,因为它“在所有输入都有可用元素时发出”,即即使源 A 中存在可用元素并且您确实想要要切换到源 A,zip 将等到 B 中存在可用元素后再发出 (ref)。

另一方面,merge 将在任何来源有可用项目时立即发出。

【讨论】:

  • 是的,这就是我在其他评论中所考虑的内容,而且确实是一个有价值的观察,zip 可能并不完全正确。我只是想知道是否有开箱即用的东西,因为多路复用似乎是相当常见的操作。另外,您能否详细说明为什么使用via(flow) 而不是直接映射到Source
  • via(Flow) 确实很傻 - 删除它!我对Akka没有经验,另一种选择似乎是Source.combine,但它看起来更复杂。
【解决方案2】:

我不确定这是您需要的,但您可以尝试以下方法:

def cond: Boolean = Random.nextBoolean()

val src1 = Source.fromIterator(() => LazyList.from(1).iterator)
val src2 = Source.fromIterator(() => LazyList.from(-1, -1).iterator)
val src3 = src1.zip(src2).map(pair => if (cond) pair._1 else pair._2)

src3.runForeach(println)

在一个run in Scastie中,输出开始于:

1
-2
3
4
-5
-6
7
-8
-9
10
...

如您所见,在此示例中,我在 2 个流之间随机选择。

【讨论】:

    猜你喜欢
    • 2019-06-25
    • 1970-01-01
    • 2016-02-22
    • 1970-01-01
    • 1970-01-01
    • 2016-08-22
    • 1970-01-01
    • 2018-12-07
    • 1970-01-01
    相关资源
    最近更新 更多