【问题标题】:Limit amount of elements a stream is processing at one time限制流一次处理的元素数量
【发布时间】:2020-03-21 06:43:08
【问题描述】:

AKKA 中的一个人如何限制当前存在于(部分)流中的元素数量,而不必删除任何元素?

【问题讨论】:

  • 在堆栈溢出时感谢辛勤工作。请提供一些您卡住且无法继续前进的代码。

标签: scala akka akka-stream


【解决方案1】:

您所说的称为溢出策略。在documentation you link to 第一个例子显示了你想要的溢出策略:OverflowStrategy.backpressure

Akka 流是反应式的,这意味着通过将溢出策略设置为背压,您就是在告诉您连接的生产者您现在不能再消费任何数据。但是,这只有在您上游的所有内容都可以处理您不再处理任何元素的事实时才有效。

在你的情况下,这样的事情应该可以工作,因为我们的Source 可以被背压:

import scala.concurrent.duration._

implicit val system = ActorSystem()

FileIO.fromPath(Paths.get(...))
  .via(Compression.gunzip())
  .via(Framing.delimiter(ByteString("\n"), maximumFrameLength = 255))
  .map(_.utf8String)
  .buffer(10, OverflowStrategy.backpressure)
  .throttle(elements = 1, per = 1.second)
  .to(Sink.foreach(println))
  .run()

通过这个非常简单的示例,很容易看出背压的作用。由于 w 使用throttle,因此流量将对上游产生反压,因此每秒仅发射 1 个元素。给定一个足够大的输入文件,这应该会很快填满 10 个元素的缓冲区。

如果你换成使用OverflowStrategy.fail 并再次尝试运行它,你会发现流几乎立即失败,因为缓冲区已满:

[ERROR] [Buffer(akka://default)] Failing because buffer is full and overflowStrategy is: [Fail]

【讨论】:

    【解决方案2】:

    对于 Akka Streams 中的此类事情,我会使用 mapAsync 并注意不要使用 .async(因为后者,至少没有自定义 ActorMaterializer 或调度任务的调度程序,不会传播背压那么紧)。

    object AkkaStreamLimitInflight {
      implicit val actorSystem = ActorSystem("foo")
      implicit val mat = ActorMaterializer()
    
      def main(args: Array[String]): Unit = {
        import actorSystem.dispatcher
    
        val inflight = new AtomicInteger(0)
    
        def printWithInflight(msg: String): Unit = {
          println(s"$msg (${inflight.get} inflight)")
        }
    
        val source = Source.unfold(0) { state => 
          println(s"Emitting $state (${inflight.incrementAndGet()} inflight)")
          Some((state + 1, state))
        }.take(10)
    
        def quintuple(i: Int): (Int, Int) = {
          val quintupled = 5 * i
          printWithInflight(s"$i quintupled is $quintupled (originally $i)")
          (i, quintupled)
        }
    
        def minusOne(tup: (Int, Int)): (Int, Int) = {
          val (original, i) = tup
          val minus1 = i - 1
          printWithInflight(s"$i minus one is $minus1 (originally $original)")
          (original, minus1)
        }
    
        def double(tup: (Int, Int)): (Int, Int) = {
          val (original, i) = tup
          val doubled = 2 * i
          printWithInflight(s"$i doubled is $doubled (originally $original)")
          (original, doubled)
        }
    
        val toUnit = Flow[(Int, Int)]
          .map { case (original, i) =>
            println(s"Done with $i (originally $original) 
    (${inflight.decrementAndGet()} inflight)")
            ()
          }
    
        val fut: Future[Done] = source
          .mapAsync(1) { i => Future { quintuple(i) }}
          .mapAsync(1) { tup => Future { minusOne(tup) }.map(double) }
          .via(toUnit)
          .runWith(Sink.ignore)
    
    
        fut.onComplete { _ => actorSystem.terminate() }
      }
    }
    

    在这种情况下,飞行计数永远不会超过 4(第一个 mapAsync 之前的一个(由于操作员融合),第一个 mapAsync 一个,第二个 mapAsync 一个,第二个之后的一个mapAsync)。如果您想在特定阶段限制飞行中元素的数量,这是可行的方法。

    但是,如果您只想限制流中正在运行的工作,则将业务逻辑移动到单个 Future 并在 mapAsync 中生成有限数量的期货,只需一个旋钮即可:

    val fut: Future[Done] = source
      .mapAsync(5) { i =>
        Future { quintuple(i) }
          .map(minusOne)
          .map(double)
      }
      .via(toUnit)
      .runWith(Sink.ignore)
    

    【讨论】:

    • 感谢您的快速回复!我目前有以下内容:source.via(parsing).via(component1).to(component2).run()。在我的流的解析阶段,我不希望对通过它的元素数量有任何限制。从组件 1 的入口开始一直到水槽,对于整个流的整个部分,应该限制最多 100 个元素在飞行中。我还可能提到,component2 是一个非常复杂的组件,具有分区、广播、自定义合并和多个接收器。不知道会不会有什么影响。你的方法适用于这种描述吗?
    猜你喜欢
    • 2016-01-13
    • 1970-01-01
    • 1970-01-01
    • 2020-03-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-05-17
    相关资源
    最近更新 更多