【问题标题】:How to stop runnable graph如何停止可运行图
【发布时间】:2018-02-13 05:49:51
【问题描述】:

使用 akka 流开始我的第一步。我有一个类似于从here 复制的图表:

val topHeadSink = Sink.head[Int]
val bottomHeadSink = Sink.head[Int]
val sharedDoubler = Flow[Int].map(_ * 2)    
val g =    RunnableGraph.fromGraph(GraphDSL.create(topHeadSink, bottomHeadSink)((_, _)) { implicit builder =>
      (topHS, bottomHS) =>
      import GraphDSL.Implicits._
      val broadcast = builder.add(Broadcast[Int](2))
      Source.single(1) ~> broadcast.in       

  broadcast.out(0) ~> sharedDoubler ~> topHS.in
  broadcast.out(1) ~> sharedDoubler ~> bottomHS.in
  ClosedShape
})

我可以使用g.run() 运行图表 但我怎么能阻止它? 在什么情况下我应该这样做(除了不使用 - 商业方面)? 该图包含在一个参与者中。如果 Actor 崩溃了,actor 底层的图会发生什么?它也会终止吗?

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    documentation 中所述,从图形外部完成图形的方法是使用KillSwitch。您从文档中复制的示例不能很好地说明这种方法,因为源只是一个元素,当您运行它时,流将很快完成。让我们调整图表以更轻松地查看 KillSwitch 的实际效果:

    val topSink = Sink.foreach(println)
    val bottomSink = Sink.foreach(println)
    val sharedDoubler = Flow[Int].map(_ * 2)
    val killSwitch = KillSwitches.single[Int]
    
    val g = RunnableGraph.fromGraph(GraphDSL.create(topSink, bottomSink, killSwitch)((_, _, _)) {
      implicit builder => (topS, bottomS, switch) =>
    
      import GraphDSL.Implicits._
    
      val broadcast = builder.add(Broadcast[Int](2))
      Source.fromIterator(() => (1 to 1000000).iterator) ~> switch ~> broadcast.in
    
      broadcast.out(0) ~> sharedDoubler ~> topS.in
      broadcast.out(1) ~> sharedDoubler ~> bottomS.in
      ClosedShape
    })
    
    val res = g.run // res is of type (Future[Done], Future[Done], UniqueKillSwitch)
    Thread.sleep(1000)
    res._3.shutdown()
    

    源现在由一百万个元素组成,接收器现在打印广播的元素。在我们调用shutdown 来完成流之前,流会运行一秒钟,这不足以搅动所有一百万个元素。

    如果您在 Actor 内部运行流,则为运行流而创建的底层 Actor(或多个 Actor)的生命周期是否与“封闭”Actor 的生命周期相关,这取决于物化器的创建方式。阅读documentation 了解更多信息。以下由 Colin Breck 撰写的关于使用 actor 和 KillSwitch 来管理流的生命周期的博文也很有帮助:http://blog.colinbreck.com/integrating-akka-streams-and-akka-actors-part-ii/

    【讨论】:

      【解决方案2】:

      有一个适合您的 KillSwitch 功能。检查其他 SO 问题的答案:Proper way to stop Akka Streams on condition

      【讨论】:

        猜你喜欢
        • 2014-11-04
        • 2011-08-16
        • 1970-01-01
        • 2018-09-26
        • 1970-01-01
        • 2013-11-15
        • 2019-10-13
        • 2013-03-08
        • 1970-01-01
        相关资源
        最近更新 更多