【问题标题】:Akka streams dynamic filtersAkka 流式动态过滤器
【发布时间】: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 =&gt; filterNames.exists(_.equalsIgnoreCase(person.name)))).to(printSink).run()
  • 对于广播到合并,请注意,您将获得每个元素的 n 个发射。这是故意的吗?
  • @LeviRamsey 是的,对于每个发射,我计划对每个过滤的元素流应用一个函数。
  • 该功能会在合并之后吗?我指出合并将发出每个传入元素,无论它通过过滤器多少次。
  • @invzbl3 是的,没错。

标签: scala akka-stream


【解决方案1】:

在伪代码中,它看起来像这样,但这是一个很好的选择:

def filterAll[A](stream: AkkaStrem[A])(filters: List[A => Boolean]): AkkaStrem[A] =
  stream.filter(a => filters.forall(p => p(a)))

要点和技巧是使用 forall 进行单个过滤器。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-11-27
    • 2019-11-16
    • 2018-12-17
    • 1970-01-01
    • 2018-01-10
    • 1970-01-01
    • 1970-01-01
    • 2016-06-20
    相关资源
    最近更新 更多