我有一些類似于以下的代碼:
object Test extends App {
val SomeComplicatedFlow: Flow[Int, Int, NotUsed] =
Flow.fromGraph(GraphDSL.create() { implicit builder =>
import GraphDSL.Implicits._
val input = builder.add(Balance[Int](1)) //Question 1) how to get rid of this input
val buffer = Flow[Int].buffer(12, OverflowStrategy.backpressure)
val balance = builder.add(Balance[Int](2))
val flow1 = Flow[Int].map(_*2)
val flow2 = Flow[Int].map(_*2)
val zip = builder.add(ZipWith[Int, Int, Int]((left, right) => {
left right
}))
val flow3 = Flow[Int].map(_*2)
input ~> buffer ~> balance.in
balance.out(0) ~> flow1 ~> zip.in0
balance.out(1) ~> flow2 ~> zip.in1
zip.out ~> flow3
FlowShape(input.in, flow3) //Question 2) how to make an outlet here
})
}
請注意,我必須添加一個Balance被呼叫的input,因為我無法Inlet從我想要創建的第一個Buffer中檢索一個FlowShape。有沒有其他更簡單的方法來解決這個問題?創建一個Balancewith 1Outlet似乎是錯誤的方法。
我的第二個問題是類似的。我無法檢索Outletfrom flow3。我知道解決這個問題的唯一方法是創建另一個Balance,并將其Outlet作為Outlet整個FlowShape. 有沒有更好的方法來解決這個問題?
uj5u.com熱心網友回復:
ABalance是向第一個可用輸出發出的扇出形狀。考慮到您要在下一步中壓縮流,您需要的是Broadcast. 當所有輸出都可用時,它將扇出到所有輸出。
此外,構建器可以添加任何形狀,Graph包括Flow。您不必為此使用自定義形狀。
更新后的代碼:
object Test extends App {
val SomeComplicatedFlow: Flow[Int, Int, NotUsed] =
Flow.fromGraph(GraphDSL.create() { implicit builder =>
import GraphDSL.Implicits._
val buffer = Flow[Int].buffer(12, OverflowStrategy.backpressure)
val input = builder.add(buffer)
val broadcast = builder.add(Broadcast[Int](2))
val flow1 = Flow[Int].map(_*2)
val flow2 = Flow[Int].map(_*2)
val zip = builder.add(ZipWith[Int, Int, Int]((left, right) => {
left right
}))
val flow3 = builder.add(Flow[Int].map(_*2))
input ~> broadcast.in
broadcast.out(0) ~> flow1 ~> zip.in0
broadcast.out(1) ~> flow2 ~> zip.in1
zip.out ~> flow3.in
FlowShape(input.in, flow3.out)
})
}
轉載請註明出處,本文鏈接:https://www.uj5u.com/qukuanlian/368361.html
