【问题标题】:Akka HTTP Stream listener stops processing databytes after a whileAkka HTTP Stream 侦听器在一段时间后停止处理数据字节
【发布时间】: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


【解决方案1】:

这就是原因:

EntityStreamSizeException:实际实体大小(无)超出内容长度限制(8388608 字节)!您可以通过设置 akka.http.[server|client].parsing.max-content-length 或在具体化 dataBytes 流之前调用 HttpEntity.withSizeLimit 来配置它。

对于在连续响应流的情况下寻求解决方案的任何人,您都可以通过这种方式获取源,使用 withoutSizeLimit:

val source: Source[ByteString, Any] = response.entity.withoutSizeLimit().dataBytes

【讨论】:

    猜你喜欢
    • 2019-10-20
    • 2016-12-29
    • 2017-07-28
    • 2016-06-17
    • 1970-01-01
    • 1970-01-01
    • 2020-03-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多