【发布时间】:2022-10-24 00:13:08
【问题描述】:
我正在尝试创建过滤器的动态列表,因为我需要过滤 100 个项目并且为每个项目应用一个函数,我不想为每个过滤器显式定义一个出口,因此定义了动态过滤器:
import akka.NotUsed
import akka.actor.ActorSystem
import akka.stream.ClosedShape
import akka.stream.scaladsl.{Broadcast, Flow, GraphDSL, Merge, RunnableGraph, Sink, Source}
object DynamicFilters extends App {
implicit val actorSystem = ActorSystem()
case class Person(name: String, age: Double)
val filterNames = List("1" , "2" , "3");
val printSink = Sink.foreach[Person](println)
val input = Source(List(Person("1", 30),Person("1", 20),Person("1", 20),Person("1", 30),Person("2", 2)))
val graph = RunnableGraph.fromGraph(
GraphDSL.create() { implicit builder: GraphDSL.Builder[NotUsed] =>
import GraphDSL.Implicits._
val broadcast = builder.add(Broadcast[Person](filterNames.size))
val merge = builder.add(Merge[Person](filterNames.size))
input ~> broadcast
for(index <- 0 to filterNames.size-1){
println("Adding filter")
val fi = Flow[Person].filter(f => f.name.equalsIgnoreCase(filterNames(index)))
broadcast.out(index) ~> fi ~> merge
}
merge ~> printSink
ClosedShape
}
)
graph.run()
}
这个解决方案看起来很“hacky”,是否有一种替代方法使用 Akka 流来过滤图中的许多项目而不为每个项目定义自定义出口?
【问题讨论】:
-
为什么不
input.via(Flow[Person].filter(person => filterNames.exists(_.equalsIgnoreCase(person.name)))).to(printSink).run()? -
对于广播到合并,请注意,您将获得每个元素的 n 个发射。这是故意的吗?
-
@LeviRamsey 是的,对于每个发射,我计划对每个过滤的元素流应用一个函数。
-
该功能会在合并之后吗?我指出合并将发出每个传入元素,无论它通过过滤器多少次。
-
@invzbl3 是的,没错。
标签: scala akka-stream