【问题标题】:Is there a higher-level way to write a custom GraphStage?是否有更高级别的方法来编写自定义 GraphStage?
【发布时间】:2016-02-04 20:50:50
【问题描述】:

我正在尝试了解 Akka Streams 的一些较新的部分。我有一个自定义 FanOutShape2 形状,它做一些非常简单的事情:输入 (Boolean, Option[A_Thing]) 并决定是否将流路由到 out0 或 out1(通过或失败),如下所示:

object PassFilter {
  type FilterShape = FanOutShape2[(Boolean, Option[OutputWrapper]), OutputWrapper, akka.NotUsed]
}
import PassFilter._

case class PassFilter()(implicit asys: ActorSystem) extends GraphStage[FilterShape] {
  val mergedIn: Inlet[(Boolean, Option[OutputWrapper])] = Inlet("Merged")
  val outPass: Outlet[OutputWrapper] = Outlet("Pass")
  val outFail: Outlet[akka.NotUsed] = Outlet("Fail")

  override val shape: FilterShape = new FanOutShape2(mergedIn, outPass, outFail)

  override def createLogic(inheritedAttributes: Attributes): GraphStageLogic = new GraphStageLogic(shape) {
    override def preStart(): Unit = pull(mergedIn)
    setHandler(mergedIn, new InHandler {
      override def onPush(): Unit = {
        val (passedPrivacy, outWrapper) = grab(mergedIn)
        if (!passedPrivacy || outWrapper.isEmpty)
          push(outFail, akka.NotUsed)
        else
          push(outPass, outWrapper.get)
        pull(mergedIn)
      }
      override def onUpstreamFinish(): Unit = {} // necessary for some reason!
    })
    setHandler(outPass, eagerTerminateOutput)
    setHandler(outFail, eagerTerminateOutput)
  }
}

这个基本的想法是可行的,我可以看到它给了我完全控制过程的可能性,但是对于这个超级琐碎的决策逻辑,有一个更简单、更高级别的方法来创建一个“小部件”,可以包含在我的 GraphDSL 中?

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    此答案基于akka-stream 版本2.4.2-RC1。 API 在其他版本中可能略有不同。依赖可以被sbt消费:

    libraryDependencies += "com.typesafe.akka" %% "akka-stream" % "2.4.2-RC1"
    

    您可以轻松创建自己的逻辑抽象。您可以将其添加为类型参数,而不是硬编码 PassFilter 类的类型。您可以将其作为构造函数参数传递,而不是硬编码确定输出端口的函数。通过这样做,您将收到一个可重用的组件,您可以连接到任意流。幸运的是,Akka 已经提供了这样一个组件。它被称为Partition

    val shape = GraphDSL.create() { implicit b ⇒
      import GraphDSL.Implicits._
    
      val first = b.add(Sink.foreach[Int](elem ⇒ println("even:\t" + elem)))
      val second = b.add(Sink.foreach[Int](elem ⇒ println("odd:\t" + elem)))
      val p = b.add(Partition[Int](2, elem ⇒ if (elem%2 == 0) 0 else 1))
    
      p ~> first
      p ~> second
    
      SinkShape(p.in)
    }
    Source(1 to 5).to(shape).run()
    
    /*
    This should print:
    odd:    1
    even:   2
    odd:    3
    even:   4
    odd:    5
    */
    

    Partition 组件还采用多个输出端口作为参数,以使其更加可重用。使用~> 符号,您可以将输出端口连接到其他组件,正如我在示例中所做的那样。您当然可以不考虑它并通过FanOutShape2 组件返回两个输出端口。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-05-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-06-13
      • 2017-11-30
      • 2022-12-31
      • 1970-01-01
      相关资源
      最近更新 更多