【发布时间】:2019-01-09 20:59:15
【问题描述】:
在服务器端,我正在使用一个 HTTP API,它在页面中返回其结果。如,响应包含 x 数量的结果,如果超过 0,我可以再次调用它以获得下一个 x 结果。 x 可以任意选择,直到 API 的最大页面大小。
现在我想通过 WebSocket 高效地流式传输所有结果集,而不会压倒它(应用背压)。最初我构建了整个结果集,然后从中创建了一个 Source:
getEventsFuture().foreach { events =>
sender ! Flow.fromSinkAndSource(Sink.ignore, Source(events))
}
这有效,WebSocket 客户端以最大速度接收所有事件。这样做的最大缺点是我必须在开始将数据返回给我的客户端之前获取所有页面。理想情况下,我会使用较小的页面大小,并在客户端连接后立即将结果返回给客户端,并在此过程中获取下一页。
因此,我需要一个带有源的流,我可以在流实现后向其中添加数据。我为此尝试使用Source.actorRef:
val events = Source.actorRef[Event](1000, OverflowStrategy.fail).mapMaterializedValue { outActor =>
sendEvents(outActor)
NotUsed
}
sender ! Flow.fromSinkAndSource(Sink.ignore, events)
基本上,我采用物化的 actorRef 并将所有事件发送给它。每次获取页面时,我都会将结果转储给演员。现在,我对 Source 的初始化可能已经告诉您这并不总是有效。有时,当响应足够大并且客户端没有像其他时候那样快速消费时,套接字连接就会关闭。我觉得OverflowStrategy.fail 是反对丢弃事件的正确策略,因为我不希望客户认为他们得到了一切,如果不是这样的话。
我没有预先为缓冲区设置合理的值,我不想设置 Int.max 或其他东西,因为我认为 Akka 内部确实为缓冲区大小分配了全部内存。
我该如何解决这个问题?我希望所有事件都尽可能快地发送给客户端,并且像第一个示例一样具有适当的背压。
获取第一页后,我确实知道总共会有多少结果,因此我可以预先获取一个小页面并将缓冲区大小设置为完整结果大小,但这似乎是一种解决方法。
【问题讨论】:
-
我仍在尝试完全理解上下文,但是如果您使用
Source.fromIterator会怎样?您的迭代器可以获取每个页面并将其发送到连接的客户端。关于由于处理时间而关闭的套接字,您可能想尝试使用类似.keepAlive的空字节字符串。
标签: scala akka akka-stream akka-http