【问题标题】:How to work Supervision strategy after recover at non-linear graphs in Akka Streams在 Akka Streams 中恢复非线性图后如何工作监督策略
【发布时间】: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


    【解决方案1】:

    首先,让我们添加几个打印语句以更清楚地了解发生了什么:一个在流完成时...

    val stream =
      Source(1 to 10)
        .via(graph)
        .withAttributes(ActorAttributes.supervisionStrategy(Supervision.resumingDecider))
        .runForeach(println)
    
    // ...
    
    stream.onComplete { _ =>
      println("Done!") // <---
      system.terminate()
    }
    

    ...recover 块中的另一个:

    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 =>
            println("Recovering...") // <---
            Failure[String](e)
        }
    }
    

    运行流的输出是...

    (Success(1),hello1)
    (Success(2),hello2)
    (Success(3),hello3)
    (Success(4),hello4)
    // no "Recovering..." or "Done!"
    

    ...表明recover 方法没有被调用并且流永远不会完成。流死锁的原因与blog 描述的原因相同:

    • [dangerFlow] 失败并且不会向Zip 发出元素。然后它继续从broadcast 请求下一个元素。但是,broadcast 要发出元素,必须从所有输出发出需求信号。

    • Zip 只接收一个元素(来自safeFlow)并永远等待第二个元素。 Zip 仅在两个输入都有值时发出。

    恢复监管策略是recover未被调用的原因。删除该策略...

    val stream =
      Source(1 to 10)
        .via(graph)
        //.withAttributes(ActorAttributes.supervisionStrategy(Supervision.resumingDecider))
        .runForeach(println)
    

    ...再次运行流会产生以下输出:

    (Success(1),hello1)
    (Success(2),hello2)
    (Success(3),hello3)
    (Success(4),hello4)
    Recovering...
    (Failure(java.lang.Exception: 5 is invalid),hello5)
    Done!
    

    现在recover 被调用,流完成,但流被截断。这是因为recover 完成了流:

    recover 允许您发出最终元素,然后在上游失败时完成流。

    要获得所需的行为,您必须使用Try,如下所示:

    private def dangerFlow: Flow[Int, Try[String], NotUsed] = {
      Flow[Int].map(a => if (a == 5) Failure(new Exception("5 is invalid")) else Try(a.toString))
    }
    

    使用上述Flow 运行流会产生以下结果:

    (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)
    Done!
    

    【讨论】:

    • 感谢您的回复!我理解recover“完成”流,所以如果监督策略正在恢复,recover 不会被调用。我必须使用Try 而不是传播异常。
    猜你喜欢
    • 1970-01-01
    • 2015-11-14
    • 2018-03-17
    • 2019-03-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-16
    • 1970-01-01
    相关资源
    最近更新 更多