【发布时间】:2018-06-20 14:50:50
【问题描述】:
是否有与mapAsync() 方法等效的方法,但对于filter?
这是一个使用伪代码的例子:
val filter: T => Future[Boolean] = /.../
source.filter(filter).runWith(/.../)
^^^^^^
谢谢
【问题讨论】:
标签: scala akka akka-stream
是否有与mapAsync() 方法等效的方法,但对于filter?
这是一个使用伪代码的例子:
val filter: T => Future[Boolean] = /.../
source.filter(filter).runWith(/.../)
^^^^^^
谢谢
【问题讨论】:
标签: scala akka akka-stream
我不认为Flow 或Source 的直接方法具有您正在寻找的功能,但可用方法的组合将得到您想要的:
def asyncFilter[T](filter: T => Future[Boolean], parallelism : Int = 1)
(implicit ec : ExecutionContext) : Flow[T, T, _] =
Flow[T].mapAsync(parallelism)(t => filter(t).map(_ -> t))
.filter(_._1)
.map(_._2)
【讨论】:
(Boolean, T) 元组作为内部map 阶段的输出,对吧? Flow[T].mapAsync(parallelism)(t => filter(t).map(b -> (b, t))).filter(_._1).map(_._2)
_ -> t 确实构造了一个元组。在您的示例中,您需要.map(b => (b, t);我写它的方式是等效的简写。注意:-> 创建一个元组,=> 创建一个函数。