【问题标题】:How to generate a big data stream on the fly如何动态生成大数据流
【发布时间】:2014-05-19 16:40:11
【问题描述】:

我必须即时生成一个大文件。读取数据库并将其发送给客户端。 我阅读了一些文档并做到了这一点

val streamContent: Enumerator[Array[Byte]] = Enumerator.outputStream {
        os => 
              // new PrintWriter() read from database and for each record 
              // do some logic and write
              // to outputstream
      }
      Ok.stream(streamContent.andThen(Enumerator.eof)).withHeaders(
              CONTENT_DISPOSITION -> s"attachment; filename=someName.csv"
        )

我对 scala 比较陌生,一般来说只有一周时间,所以不要为我的声誉提供指导。

我的问题是:

1) 这是最好的方法吗?我发现如果我有一个大文件,这将加载到内存中,并且在这种情况下也不知道块大小是多少,如果它会为每个 write() 发送不方便。

2) 我发现这个方法Enumerator.fromStream(data : InputStream, chunkedSize : int) 稍微好一点,因为它有一个块大小,但我没有 inputStream 因为我是动态创建文件的。

【问题讨论】:

  • 块大小默认设置为 1024 * 8 。我认为大小取决于您是否要发送更大或更小的块。
  • @goral 你怎么知道Enumerator.outputStream1024*8 分块??
  • 我不知道 outputStream,你说的是 fromStream 方法和在设置为默认值的规范中
  • @goral 抱歉,如果我的问题不清楚,我已经知道 fromStream,我问的是 .outputStream

标签: scala playframework-2.1


【解决方案1】:

docs for Enumerator.outputStream中有注释:

不是 [sic!] 调用 write 不会阻塞,所以如果被馈送到的迭代器消耗输入的速度很慢,OutputStream 将不会推回。这意味着它不应该用于大型流,因为存在内存不足的风险。

是否会发生这种情况取决于您的情况。如果您可以并且将在几秒钟内生成千兆字节,您可能应该尝试不同的方法。我不确定是什么,但我将从Enumerator.generateM() 开始。但是,对于许多情况,您的方法非常好。看看at this example by Gaëtan Renaudeau for serving a Zip file that's generated on the fly in the same way you're using it:

val enumerator = Enumerator.outputStream { os =>
  val zip = new ZipOutputStream(os);
  Range(0, 100).map { i =>
    zip.putNextEntry(new ZipEntry("test-zip/README-"+i+".txt"))
    zip.write("Here are 100000 random numbers:\n".map(_.toByte).toArray)
    // Let's do 100 writes of 1'000 numbers
    Range(0, 100).map { j =>
      zip.write((Range(0, 1000).map(_=>r.nextLong).map(_.toString).mkString("\n")).map(_.toByte).toArray);
    }
    zip.closeEntry()
  }
  zip.close()
}
Ok.stream(enumerator >>> Enumerator.eof).withHeaders(
  "Content-Type"->"application/zip", 
  "Content-Disposition"->"attachment; filename=test.zip"
)

请注意,在 Play 的较新版本中,Ok.stream 已替换为 Ok.chunked,以防您想升级。

至于块大小,您始终可以使用Enumeratee.grouped 来收集一堆值并将它们作为一个块发送。

val grouper = Enumeratee.grouped(  
  Traversable.take[Array[Double]](100) &>> Iteratee.consume()  
)

然后你会做类似的事情

Ok.stream(enumerator &> grouper >>> Enumerator.eof)

【讨论】:

  • +1 感谢您的回答,非常有用,我在 scala 中非常菜鸟,我仍然很难轻松阅读和理解代码。但是在您提供的链接中,此演示展示了如何即时生成 zip 文件并将其直接流式传输到 HTTP 客户端 无需将其加载到内存中或将其存储在文件中. 你所说的千兆字节与 outOfMemory 是矛盾的
  • 是的,我知道,Scala 很难。我也花了一段时间,我还是个初学者。但是您可以在开头使用更多类似 Java/OO 的代码,然后使用 start experimenting with functional programming。不要从字面上理解“没有将其加载到内存中”。您必须将 something 存储在内存中。这只是意味着您不会在内存中生成整个 zip 文件,然后将其发送到客户端,而是在生成时将其流式传输。只要客户端出现并接收流,它就会以同样快的速度从内存中删除。
  • 好的,例如Enumerator.forStream() 让您可以设置一个chunkSize,在这种情况下,当它要分块时,就像1024*8 中的注释或输出流的每次写入一样?如果是grouper 中的魔法代码有用的话:D
  • Traversable.take 无法编译我需要导入一些特殊的东西吗?
  • 是的,导入 play.api.libs.iteratee.Traversable(如果将来遇到,可以在 Play 的 scaladocs 中搜索该类)。如果您查看Enumerator.fromStream(),您可以看到它只是从 InputStream 中读取一定数量的字节到一个数组中并将其馈送到 Enumerator 中。默认为1024 * 8。正如@goral 所暗示的那样,我不确定枚举器在发送时是否会被重新分块为 8k 块。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2010-12-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-07-06
相关资源
最近更新 更多