【问题标题】:Streaming data in and out simultaneously on a single HTTP connection in play在单个 HTTP 连接上同时输入和输出数据
【发布时间】: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


【解决方案1】:

您想要做的事情在 Play 中完全是不可能的。

问题是 Play 在完全收到请求之前无法开始发送响应。因此,您可以像以前一样接收完整的请求,然后发送响应,您可以在收到请求时处理请求(在自定义 BodyParser 中),但您仍然在您收到完整的请求之前无法回复(这是文档中的注释所暗示的 - 尽管您可以在不同的连接中发送响应)。

要了解原因,请注意Action 本质上是(RequestHeader) =&gt; Iteratee[Array[Byte], SimpleResult]。在任何时候,Iteratee 都处于以下三种状态之一 - DoneContError。它只有在Cont状态下才能接受更多的数据,但它只有在Done状态下才能返回一个值。由于该返回值为SimpleResult(即我们的响应),这意味着从接收数据到发送数据之间存在硬性中断。

根据this answer,HTTP 标准确实允许在请求完成之前做出响应,但大多数浏览器不遵守规范,并且无论如何 Play 都不支持它,如上所述。

在 Play 中实现全双工通信的最简单方法是使用 WebSocket,但我们已经排除了这种可能性。如果服务器资源使用是更改的主要原因,您可以尝试使用play.api.mvc.BodyParsers.parse.temporaryFile 解析数据,这会将数据保存到临时文件,或者play.api.mvc.BodyParsers.parse.rawBuffer,如果请求过多,则会溢出到临时文件大。

否则,我看不到使用 Play 执行此操作的合理方法,因此您可能需要考虑使用其他网络服务器。

【讨论】:

  • "...HTTP 标准确实允许在请求完成之前做出响应,但大多数浏览器不遵守规范..." 因为客户端不是在我的情况下是浏览器,这不是问题。 “无论如何 Play 都不支持它。” - 你确定吗?我找不到任何相关的信息。无论如何,我在想的是,在回复之前不要让Iteratee 达到它的Done 状态,并且在解析数据时以某种方式使用连接开始流式传输数据。如果不可能,我将不胜感激。谢谢。
  • 这是不可能的,因为一旦Iteratee处于Done状态,就会产生响应。如果您查看 Iteratee 的 API,当它处于任何其他状态时,无法从中获得结果。状态由 3 个案例类 Cont、Done 和 Err 表示,并且只有 Done 包含结果。这是产生响应的唯一方式。
  • 此外,如果您查看 Play 的 Web 服务器的实现(在 play.core 类中),您会发现一旦 Iteratee 处于 Done 状态,它就会停止处理数据,所以如果您正在考虑以某种方式欺骗它(比如在你的迭代中覆盖折叠,同时调用 Done 和 Cont 回调),这也行不通。
  • 既然我想做一些“非正统”的事情,我可以有一个“hacky”的解决方案,而不是走“标准”的方式。如果有可能以某种方式在BodyParser 中获得连接并开始发送响应,我会接受的。另一个“hacky”选项是以某种方式将Iteratee 留在Cont 模式下,因此它仍会消耗传入的数据,但会急切地调用该操作,Enumerator 发出解析的块,可用作请求正文。
  • 基本问题是,唯一您可以向 Play 提供 Enumerator 的方法是让您的 Iteratee 返回一个 SimpleResult,然后 发生这种情况的唯一方法是您的Iteratee 处于Done 状态。我确实考虑了一个像你提议的那样的 hacky 选项 - Iteratee 通过调用三个回调之一来表示其状态,所以我想知道是否可以通过调用 Done 和来打破 Iteratee 合同Cont 回调。不幸的是,由于我之前评论中概述的原因,这也不起作用。
【解决方案2】:

“在单个 HTTP 连接上同时输入和输出数据流”

我还没有读完你的所有问题,也没有读完代码,但是你要求做的事情在 HTTP 中不可用。这与 Play 无关。

当您发出网络请求时,您打开一个到网络服务器的套接字并发送“GET /file.html HTTP/1.1\n[可选标头]\n[更多标头]\n\n”

在您完成请求后(并且仅在此之后)您会收到回复(可选地包括请求正文作为请求的一部分)。当且仅当请求 响应完成时,在 HTTP 1.1(但不是 1.0)中,您可以在同一个套接字上发出新请求(在 http 1.0 中,您打开一个新套接字)。

响应“挂起”是可能的……这就是网络聊天的工作方式。服务器只是坐在那里,挂在打开的套接字上,在有人向您发送消息之前不发送响应。当/如果您收到聊天消息时,与 Web 服务器的持久连接最终会提供响应。

同样,请求也可以“挂起”。您可以开始将您的请求数据发送到服务器,稍等片刻,然后当您收到额外的用户输入时完成请求。与在每个用户输入上不断创建新的 http 请求相比,这种机制提供了更好的性能。服务器可以将此数据流解释为不同输入的流,即使这不一定是 HTTP 规范的最初意图。

HTTP 不支持接收部分请求,然后发送部分响应,然后接收更多请求的机制。它只是不在规范中。一旦您开始接收响应,向服务器发送附加信息的唯一方法是使用另一个 HTTP 请求。您可以使用已经并行打开的一个,也可以打开一个新的,或者您可以完成第一个请求/响应并在同一个套接字上发出另一个请求(在 1.1 中)。

如果您必须在单个套接字连接上使用异步 io,您可能需要考虑使用 HTTP 以外的其他协议。

【讨论】:

  • 从技术上讲,您可以在收到响应之前发送另一个请求,就像在请求管道中一样
  • 有趣,我不知道请求管道。根据维基百科“在所有其他浏览器 [除了歌剧] HTTP 管道被禁用或未实现。”它似乎仍然不会启用请求的功能。 “HTTP 1.1 的限制仍然适用:服务器必须按照接收请求的顺序发送响应……连接保持先进先出。” OP“在这种情况下,只有在请求完成后,我才会从请求中获取正文,并且所有数据都已上传到服务器。但我想随时开始为请求提供服务。”
  • 是的,请求流水线在很大程度上与 OP 试图做的事情无关。请求管道与此相反,它允许客户端在收到回复之前发送请求,而不是服务器在收到完整请求之前发送回复。
猜你喜欢
  • 1970-01-01
  • 2019-07-18
  • 1970-01-01
  • 1970-01-01
  • 2017-03-10
  • 2013-06-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多