【发布时间】: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