You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 10:14:40