【发布时间】:2017-10-31 09:51:58
【问题描述】:
我有一个包含广播和压缩的流程图。如果在这个流程中有什么(不管它是什么)失败,我想删除传递给它的有问题的元素并恢复。我想出了以下解决方案:
val flow = Flow.fromGraph(GraphDSL.create() { implicit builder =>
import GraphDSL.Implicits._
val dangerousFlow = Flow[Int].map {
case 5 => throw new RuntimeException("BOOM!")
case x => x
}
val safeFlow = Flow[Int]
val bcast = builder.add(Broadcast[Int](2))
val zip = builder.add(Zip[Int, Int])
bcast ~> dangerousFlow ~> zip.in0
bcast ~> safeFlow ~> zip.in1
FlowShape(bcast.in, zip.out)
})
Source(1 to 9)
.via(flow)
.withAttributes(ActorAttributes.supervisionStrategy(Supervision.restartingDecider))
.runWith(Sink.foreach(println))
我希望它打印出来:
(1,1)
(2,2)
(3,3)
(4,4)
(5,5)
(6,6)
(7,7)
(8,8)
(9,9)
但是,它会死锁,只打印:
(1,1)
(2,2)
(3,3)
(4,4)
我们已经进行了一些调试,结果发现它对其子代应用了“恢复”策略,这导致dangerousFlow 在失败后恢复,从而要求bcast 提供一个元素。 bcast 不会发出元素,直到 safeFlow 需要另一个元素,这实际上从未发生过(因为它正在等待来自 zip 的请求)。
有没有办法恢复图表而不管其中一个阶段出现了什么问题?
【问题讨论】:
标签: scala stream akka akka-stream