【发布时间】:2013-12-11 15:02:28
【问题描述】:
流式传输数据非常容易。
这是我打算如何做的一个简单示例(如果我做错了,请告诉我):
def getRandomStream = Action { implicit req =>
import scala.util.Random
import scala.concurrent.{blocking, ExecutionContext}
import ExecutionContext.Implicits.global
def getSomeRandomFutures: List[Future[String]] = {
for {
i <- (1 to 10).toList
r = Random.nextInt(30000)
} yield Future {
blocking {
Thread.sleep(r)
}
s"after $r ms. index: $i.\n"
}
}
val enumerator = Concurrent.unicast[Array[Byte]] {
(channel: Concurrent.Channel[Array[Byte]]) => {
getSomeRandomFutures.foreach {
_.onComplete {
case Success(x: String) => channel.push(x.getBytes("utf-8"))
case Failure(t) => channel.push(t.getMessage)
}
}
//following future will close the connection
Future {
blocking {
Thread.sleep(30000)
}
}.onComplete {
case Success(_) => channel.eofAndEnd()
case Failure(t) => channel.end(t)
}
}
}
new Status(200).chunked(enumerator).as("text/plain;charset=UTF-8")
}
现在,如果您通过此操作获得服务,您将获得如下信息:
after 1757 ms. index: 10.
after 3772 ms. index: 3.
after 4282 ms. index: 6.
after 4788 ms. index: 8.
after 10842 ms. index: 7.
after 12225 ms. index: 4.
after 14085 ms. index: 9.
after 17110 ms. index: 1.
after 21213 ms. index: 2.
after 21516 ms. index: 5.
在随机时间过去后接收每一行。
现在,假设我想在将数据从服务器流式传输到客户端时保留这个简单的示例,但我还想支持从客户端到服务器的完整数据流式传输。
所以,假设我正在实现一个新的BodyParser,它将输入解析为List[Future[String]]。这意味着,现在,我的 Action 可能看起来像这样:
def getParsedStream = Action(myBodyParser) { implicit req =>
val xs: List[Future[String]] = req.body
val enumerator = Concurrent.unicast[Array[Byte]] {
(channel: Concurrent.Channel[Array[Byte]]) => {
xs.foreach {
_.onComplete {
case Success(x: String) => channel.push(x.getBytes("utf-8"))
case Failure(t) => channel.push(t.getMessage)
}
}
//again, following future will close the connection
Future.sequence(xs).onComplete {
case Success(_) => channel.eofAndEnd()
case Failure(t) => channel.end(t)
}
}
}
new Status(200).chunked(enumerator).as("text/plain;charset=UTF-8")
}
但这仍然不是我想要实现的。在这种情况下,我只会在请求完成后从请求中获取正文,并且所有数据都已上传到服务器。但我想开始服务请求。一个简单的演示,就是将收到的任何线路回显给用户,同时保持连接处于活动状态。
所以这是我目前的想法:
如果我的BodyParser 会返回Enumerator[String] 而不是List[Future[String]],会怎样?
在这种情况下,我可以简单地执行以下操作:
def getParsedStream = Action(myBodyParser) { implicit req =>
new Status(200).chunked(req.body).as("text/plain;charset=UTF-8")
}
所以现在,我面临如何实现这样的BodyParser 的问题。
更准确地了解我到底需要什么,嗯:
我需要接收大块数据以解析为字符串,其中每个字符串都以换行符\n 结尾(尽管可能包含多行......)。每个 “行块” 都将由一些(与此问题无关的)计算处理,这将产生 String,或者更好的是 Future[String],因为此计算可能需要一些时间。此计算的结果字符串应在准备好时发送给用户,就像上面的随机示例一样。这应该在发送更多数据的同时发生。
我已经研究了几种试图实现它的资源,但到目前为止都没有成功。
例如scalaQuery play iteratees -> 看起来这个人正在做与我想做的事情相似的事情,但我无法将其转化为可用的示例。 (并且从 play2.0 到 play2.2 API 的差异无济于事......)
所以,总结一下:这是正确的方法吗(考虑到我不想使用WebSockets)?如果是这样,我该如何实现这样的BodyParser?
编辑:
我刚刚偶然发现了play documentation 上关于此问题的注释,上面写着:
注意:也可以实现同种直播 使用无限的 HTTP 请求以相反的方式进行通信 由接收输入数据块的自定义
BodyParser处理,但是 那要复杂得多。
所以,我不会放弃,现在我确定这是可以实现的。
【问题讨论】:
-
我不确定 HTTP 是否支持你正在尝试做的事情,看起来也不值得做 - 我读过的所有内容都表明浏览器不会处理响应,直到它们之后'已完成上传完整的请求。为什么不使用 Keep-Alive 连接将其分解为单独的请求?或使用两个单独的连接,例如一个上传数据的 AJAX 请求和一个永久的 JSONP iframe 降低数据?
-
就我而言,客户端(通常)不是网络浏览器,所以我并不关心浏览器是否支持它。无论如何,我说的是流式传输大量数据,我们可以随时处理这些数据。我永远不会希望所有数据都存储在我的服务器内存中,也不会存储在磁盘上。由于我的服务的性质,将其作为一个 HTTP 请求进行,这一点很重要。
-
这看起来像是一个 XY 问题,你真正想要做什么?
-
我将描述一项功能:客户端输入是主题网址。每行包含 1 个网址。服务器输出是 N-Triple 行。客户端发送的每个 url 都映射到描述主题的多个 n-triple。输入和输出可能包含数百万行。它必须是单个 http 连接还有其他原因。 (例如其他服务已经在使用它,我无法破坏 API...)
-
你可以保持一个连接,并使用多个请求。这确实是 HTTP 支持的唯一方式,这样做比尝试破解您自己的 HTTP 扩展要好得多。例如,您可以 RESTfully 分别为每个请求提供错误。您可以使用单个连接,不仅可以处理这些请求,还可以混合您提到的 API 支持的其他请求。使用 HTTP 2.0,可以同时处理这些请求并无序响应。您对分块请求和响应的 hack 不能做这些事情。
标签: scala playframework playframework-2.2 enumerator iterate