【发布时间】:2018-12-13 11:29:52
【问题描述】:
我使用recover 方法来捕获Akka Streams 中的错误或异常。它适用于线性图,但不适用于非线性图(例如广播、Zip)。 Graph with fan-in or fan-out 永远等待失败端口的原因,所以 Akka Streams 挂了。 解决方案在https://blog.softwaremill.com/akka-streams-pitfalls-to-avoid-part-2-f93e60746c58 的第 9 节中描述。
那篇文章使用 Try monad 并在 Flow 中捕获异常。这样可行。但是我使用恢复方法,因为我有很多流程,我想在一个地方捕获错误。
我准备了下面的例子,但是不行……
Source(1 to 10)
.via(graph)
.withAttributes(ActorAttributes.supervisionStrategy(Supervision.resumingDecider))
.runForeach(println)
private def dangerFlow: Flow[Int, Try[String], NotUsed] = {
Flow[Int].map(a => if (a == 5) throw new Exception("5 is invalid") else a.toString).map(str => Try(str)).recover {
case e => Failure[String](e)
}
}
private def safeFlow: Flow[Int, String, NotUsed] = Flow[Int].map( "hello" +_)
def graph = Flow.fromGraph(GraphDSL.create() { implicit b =>
import GraphDSL.Implicits._
val bcast = b.add(Broadcast[Int](2))
val zip = b.add(Zip[Try[String], String])
bcast.out(0) ~> dangerFlow ~> zip.in0
bcast.out(1) ~> safeFlow ~> zip.in1
FlowShape(bcast.in, zip.out)
})
结果:
(Success(1),hello1)
(Success(2),hello2)
(Success(3),hello3)
(Success(4),hello4)
我预计:
(Success(1),hello1)
(Success(2),hello2)
(Success(3),hello3)
(Success(4),hello4)
(Failure(java.lang.Exception: 5 is invalid),hello5)
(Success(6),hello6)
(Success(7),hello7)
(Success(8),hello8)
(Success(9),hello9)
(Success(10),hello10)
请告诉我任何解决方案。谢谢。
【问题讨论】:
标签: scala akka akka-stream