Scala标准库是否有将Stream[Byte]转为惰性分块Stream[InputStream]的函数
问题解答
Scala 标准库没有直接提供签名为Long => Stream[Byte] => Stream[InputStream]的现成工具函数,你可以基于标准库的惰性集合API和Java IO基础类自行实现,全程可以保证惰性求值特性,不会一次性加载全量数据到内存。
以下是符合要求的实现示例:
import java.io.InputStream import scala.collection.immutable.Stream def splitToChunkedStreams(chunkSize: Long)(byteStream: Stream[Byte]): Stream[InputStream] = { require(chunkSize > 0, "分片大小必须为正整数") // Stream的grouped为惰性实现,仅在用到对应分片时才会计算 byteStream.grouped(chunkSize.toInt).map { chunk => new InputStream { private var currentChunk: Stream[Byte] = chunk override def read(): Int = { if (currentChunk.isEmpty) -1 else { val b = currentChunk.head currentChunk = currentChunk.tail b & 0xFF } } // 批量读优化,提升IO性能 override def read(b: Array[Byte], off: Int, len: Int): Int = { if (currentChunk.isEmpty) -1 else { val takeLen = math.min(len, currentChunk.length) currentChunk.take(takeLen).copyToArray(b, off, takeLen) currentChunk = currentChunk.drop(takeLen) takeLen } } } }.toStream }
如果你使用的是Scala 3版本,原有的
Stream已被标记废弃,替换为LazyList,上述逻辑只需要把所有Stream替换为LazyList即可正常运行,惰性特性保持一致。
注意事项
- 标准库
grouped方法仅接受Int类型的分片大小,如果你的分片大小超过Int.MaxValue,需要自行实现按Long长度的惰性分片逻辑 - 每个分片的
InputStream必须按顺序读取完成后再读取下一个Stream元素,避免同时持有多个分片的字节数据 - 对接S3 Multipart Upload时要注意接口限制:除最后一个分片外,其他分片最小为5MB、最大为5GB,设置
chunkSize时要符合该规则
内容的提问来源于stack exchange,提问作者Zach LeFevre
相关产品推荐
相关产品推荐

