FS2 Stream如何实现读取InputStream直至流结束?
解决FS2从InputStream持续读取数据直到流结束的问题
嘿,我来帮你搞定这个FS2的读取问题~你现在遇到的核心问题是:你的read方法没有返回流终止的信号,而repeatEval会无脑重复调用,哪怕InputStream已经读完了。咱们一步步来解决:
第一步:修改读取方法,返回终止信号
首先,我们需要让read方法能够告诉FS2“流已经读完了”。InputStream的read方法在流结束时会返回-1,所以我们可以把返回类型改成IO[Option[Array[Byte]]]——当读取到内容时返回Some(实际读取的字节数组),流结束时返回None。
还要注意一个细节:InputStream的read(buf)可能只填满了buf的一部分(返回的字节数小于buf长度),这时候我们需要截取实际读取的部分,避免输出多余的空字节。修改后的方法如下:
import java.io.InputStream import java.nio.charset.StandardCharsets import cats.effect.IO def read(is: InputStream): IO[Option[Array[Byte]]] = IO { val buf = new Array[Byte](4096) val bytesRead = is.read(buf) if (bytesRead == -1) { None // 流结束,返回None } else { // 截取实际读取的字节,避免空字节 val actualBytes = if (bytesRead == buf.length) buf else buf.take(bytesRead) Some(actualBytes) } }
第二步:生成自动终止的Stream
有了带终止信号的read方法,我们就可以用FS2的repeatEval结合unNoneTerminate()来生成流——repeatEval会重复调用read(is),而unNoneTerminate()会在遇到None时自动停止流:
val stream = fs2.Stream.repeatEval(read(is)).unNoneTerminate()
第三步:正确管理InputStream资源
另外一个重要的点是:InputStream是需要手动关闭的,直接创建会导致资源泄漏。我们应该用Cats Effect的Resource来管理它的生命周期,确保流用完后自动关闭:
import java.io.File import java.io.FileInputStream import cats.effect.Resource // 创建InputStream的Resource:创建时打开文件,释放时关闭流 val inputStreamResource: Resource[IO, FileInputStream] = Resource.make(IO(new FileInputStream(new File("/tmp/my-file.mf"))))(is => IO(is.close()))
完整的示例代码
把这些部分组合起来,完整的代码如下:
import java.io.{File, FileInputStream, InputStream} import java.nio.charset.StandardCharsets import cats.effect.{IO, Resource} import fs2.Stream object Fs2Example { def main(args: Array[String]): Unit = { // 管理InputStream的资源 val inputStreamResource: Resource[IO, FileInputStream] = Resource.make(IO(new FileInputStream(new File("/tmp/my-file.mf"))))(is => IO(is.close())) // 使用资源并处理流 inputStreamResource.use { is => val stream = Stream.repeatEval(read(is)).unNoneTerminate() // 这里可以替换成你实际要做的处理,比如打印内容 stream.evalMap(buf => IO(println(new String(buf, StandardCharsets.UTF_8)))) .compile.drain }.unsafeRunSync() } def read(is: InputStream): IO[Option[Array[Byte]]] = IO { val buf = new Array[Byte](4096) val bytesRead = is.read(buf) if (bytesRead == -1) { None } else { val actualBytes = if (bytesRead == buf.length) buf else buf.take(bytesRead) Some(actualBytes) } } }
为什么之前的代码不行?
- 你原来的
read方法只返回字节数组,没有任何终止信号,所以repeatEval会一直调用它,哪怕InputStream已经读完了(这时候read会一直返回-1,但你还是把整个buf返回了,导致输出很多空内容)。 - 用
Option作为返回值后,unNoneTerminate()会识别到流结束的信号,自动停止流的生成。
内容的提问来源于stack exchange,提问作者St.Antario
相关产品推荐
相关产品推荐

