【发布时间】:2019-01-08 23:10:02
【问题描述】:
我们有以下架构
SQS(source) -> SQS Pollers -> 我们的业务逻辑 -> Sink 从 SQS 中删除消息。
这是一个 akka 流(我们的业务逻辑有多个阶段)。
现在我们想通过添加一个 HTTP 服务器(不是 Akka HTTP)来扩展这个架构。
现在我们的服务也有了路径
HTTP Server -> 我们的业务逻辑 -> Sink 完成一个指示 HTTP 响应完成的未来。
现在,每当 HTTP 请求到来时,我都需要一种机制来调用流。
现在 SQS 源本质上是一个长时间运行的线程,它调用服务并将消息推送到 akka 流的其余部分。
我实际上是在尝试创建一个“可调用的”akka 源,这样只有在我们收到请求时才会触发该源。
我在这里将https://doc.akka.io/docs/akka/2.5/stream/operators/Source/queue.html 视为一个潜在的解决方案,但它仅在整个可运行图实现后才返回要调用的句柄,因此合并 SQS 轮询器源和 HTTP 可调用源有点难看.
【问题讨论】:
标签: scala akka akka-stream