【问题标题】:Explicit throughput limiting on part of an akka stream对部分 akka 流的显式吞吐量限制
【发布时间】:2020-12-31 23:30:03
【问题描述】:

我的系统中有一个流程,它从 SQS 中读取一些元素(使用 alpakka)并进行一些预处理(大约 10 个阶段,通常总共不到 1 分钟)。然后,准备好的元素被发送到主处理(单个阶段,需要几分钟)。整个事情在 AWS/K8S 上运行,我们希望在 SQS 队列增长到某个阈值以上时进行扩展。问题是,SQS 队列需要很长时间才能炸毁,因为有很多元素“空闲”在进程中,已经完成了它们的预处理但正在等待主要的事情。

我们无法将预处理内容外部化到单独的队列中,因为它们的结果无法在反序列化往返中存活下来。此外,这个服务和“主”处理器是深度耦合的(这个服务作为主的 sidecar 运行)并且不能独立扩展。

从技术上讲,预处理阶段是 .mapAsyncUnordered,但整个事情已经很渺茫了(流阶段和 SQS 批处理/缓冲区)。

我们尝试降低级间缓冲区(akka.stream.materializer.max-input-buffer-size),但这只会带来一些间接的好处,没有直接的控制(而且对于我的口味来说太内部了)。

我尝试实现一个“门”包装器,它会限制任意Flow 中允许的元素数量,看起来像:

class LimitingGate[T, U](originalFlow: Flow[T, U], maxInFlight: Int) {
  private def in: InputGate[T] = ???
  private def out: OutputGate[U] = ???

  def gatedFlow: Flow[T, U, NotUsed] = Flow[T].via(in).via(originalFlow).via(out)
}

并在输入/输出门之间使用回调进行节流。

实现部分工作(流终止让我很难过),但感觉是实现实际目标的错误方法。

感谢任何想法/cmets/启发性问题

谢谢!

【问题讨论】:

  • 一个 hacky(?),但完全在流 DSL 中,限制流中飞行中的元素总数的方法,无论流中的阶段/异步边界的数量如何,它都可以工作,将使用 SQS 源压缩Source.actorRef(预实现帮助)。 zip 确保只有发送到actorRef 的每条消息都向下游传递一条消息,因此只要长时间处理完成,您就会向actorRef 发送一条消息。
  • 出于兴趣,您的代码中是否有任何asyncs(不是mapAsyncmapAsyncUnordered),因为物化器输入缓冲区只对异步边界很重要(可能在SQS 来源)。
  • @LeviRamsey ,我们的代码中没有明确的.async 边界。没有浏览所有的 alpakka 代码,但确实可能至少有一个你能解释一下“@ 987654331@ 的zip”评论吗?不知道我明白...谢谢!
  • 如果您预先实现Source.actorRef,您会得到一个ActorRef,您可以从流的其余部分向它发送消息,并将这些消息放入流中。因此,如果您有一个 zip 阶段,将 SQS 源与预先实现的 actor 源相结合,则只有在 actor 源发出一个元素时,元素才会通过。

标签: akka akka-stream


【解决方案1】:

尝试以下方法(我只是在脑海中编译它):

def inflightLimit[A, B, M](n: Int, source: Source[T, M])(businessFlow: Flow[T, B, _])(implicit materializer: Materializer): Source[B, M] = {
  require(n > 0)  // alternatively, could just result in a Source.empty...
  val actorSource = Source.actorRef[Unit](
    completionMatcher = PartialFunction.empty,
    failureMatcher = PartialFunction.empty,
    bufferSize = 2 * n,
    overflowStrategy = OverflowStrategy.dropHead  // shouldn't matter, but if the buffer fills, the effective limit will be reduced
  )
  val (flowControl, unitSource) = actorSource.preMaterialize()

  source.statefulMapConcat { () =>
    var firstElem: Boolean = true
    { a =>
      if (firstElem) {
        (0 until n).foreach(_ => flowControl.tell(()))  // prime the pump on stream materialization
        firstElem = false
      }
      List(a)
    }}
    .zip(unitSource)
    .map(_._1)
    .via(businessFlow)
    .wireTap { _ => flowControl.tell(()) }  // wireTap is Akka Streams 2.6, but can be easily replaced by a map stage which sends () to flowControl and passes through the input
 }

基本上:

  • actorSource 将为其接收到的每个 () 发出一个 Unit(),即无意义)元素
  • statefulMapConcat 将导致 n 消息仅在流首次启动时发送到 actorSource(因此允许 n 元素从源通过)
  • 仅当 actorSourcesource 都有可用的元素时,zip 才会传递来自 source() 的一对输入
  • 对于退出 businessFlow 的每个元素,将向 actorSource 发送一条消息,这将允许来自源的另一个元素通过

注意事项:

  • 这不会以任何方式限制source 内的缓冲
  • businessFlow 不能删除元素:n 元素被删除后,流将不再处理元素但不会失败;如果需要删除元素,您可以内联businessFlow 并让删除元素的阶段在删除元素时向flowControl 发送消息;你也可以做其他事情来解决这个问题

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-22
    • 1970-01-01
    • 1970-01-01
    • 2015-08-22
    • 1970-01-01
    • 2018-07-11
    相关资源
    最近更新 更多