【问题标题】:Akka streams - resuming graph with broadcast and zip after failureAkka 流 - 失败后使用广播和 zip 恢复图形
【发布时间】: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


    【解决方案1】:

    我认为你很好地理解了这个问题。您已经看到,当您的元素 5 崩溃 dangerousFlow 时,您还应该停止正在经历 safeFlow 的元素 5,因为如果它到达 zip 阶段,您就会遇到您描述的问题。我不知道如何解决您在broadcastzip 阶段之间的问题,但是将问题进一步推到哪里更容易处理呢?

    考虑使用以下dangerousFlow

    import scala.util._
    val dangerousFlow = Flow[Int].map {
      case 5 => Failure(new RuntimeException("BOOM!"))
      case x => Success(x)
    }
    

    即使出现问题,dangerousFlow 仍会发出数据。然后,您可以像目前所做的那样 zip,只需添加一个 collect 阶段作为图表的最后一步。在流程上,这看起来像:

    Flow[(Try[Int], Int)].collect {
      case (Success(s), i) => s -> i
    }
    

    现在,如您所写,如果您真的希望它输出 (5, 5) 元组,请使用以下内容:

    Flow[(Try[Int], Int)].collect {
      case (Success(s), i) => s -> i
      case (_, i)          => i -> i
    }
    

    【讨论】:

      猜你喜欢
      • 2013-12-08
      • 2013-02-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-03-19
      相关资源
      最近更新 更多