【发布时间】: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(不是mapAsync或mapAsyncUnordered),因为物化器输入缓冲区只对异步边界很重要(可能在SQS 来源)。 -
@LeviRamsey ,我们的代码中没有明确的
.async边界。没有浏览所有的 alpakka 代码,但确实可能至少有一个你能解释一下“@ 987654331@ 的zip”评论吗?不知道我明白...谢谢! -
如果您预先实现
Source.actorRef,您会得到一个ActorRef,您可以从流的其余部分向它发送消息,并将这些消息放入流中。因此,如果您有一个 zip 阶段,将 SQS 源与预先实现的 actor 源相结合,则只有在 actor 源发出一个元素时,元素才会通过。
标签: akka akka-stream