【发布时间】:2016-09-29 08:52:09
【问题描述】:
我有一个应用程序,它有 3 个像这样的 HTTP 侦听器:
val futureResponse1: Future[HttpResponse] =
Http().singleRequest(HttpRequest(uri = someUrl))
这 3 个中的每一个都在听一个不间断的流(每个都听不同的流)。并通过一个简单的流程来处理它,该流程从分组开始,然后是相对快速的处理(非阻塞):
futureResponse1.flatMap {response =>
response.status match {
case StatusCodes.OK =>
val source: Source[ByteString, Any] = response.entity.dataBytes
source.
grouped(100).
map(doSomethingFast).
runWith(Sink.ignore)
case notOK => system.log.info("failed opening, status: " + notOK.toString())
}
...
我没有收到任何异常或警告。但过了一会儿(可能是 15-25 分钟),听众突然停下来。一个接一个(不在一起)。
也许是分组阶段的问题所在?或者也许连接/流刚刚停止?或者他们共享的调度程序正在挨饿/某些东西没有被释放。
请知道为什么会发生这种情况?
==== 更新====
@Ramon J Romero 和 Vigil 我将运行更改为只有 1 个流而不是 3 个,并且删除了分组阶段。几分钟后仍然发生。我怀疑流是基于超时关闭的。我所做的就是获取块并消耗它们。
==== 更新====
找到原因,见下文。
【问题讨论】:
-
你能提供更多关于“三个人都在听不间断流”的详细信息吗?
Future值和流之间的交互可能是您问题的根源...
标签: akka akka-stream akka-http