如何在fs2流中解压.gz文件并以函数式方式逐行读取
如何用fs2函数式范式逐行读取GZIP压缩文件?
我需要逐行读取GZIP压缩文件t.gz,此前用传统IO方式实现的代码如下:
Source.fromInputStream(new GZIPInputStream(new BufferedInputStream(new FileInputStream(s))))
现在想改用fs2以函数式范式处理,已经能处理未压缩的文件,代码如下:
import cats.effect._ import fs2.io.file.{Files, Path} import fs2.{Stream, text, compression, io} object Main extends IOApp.Simple { def doThing(inPath: Path): Stream[IO, Unit] = { Files[IO] .readAll(inPath) .through(text.utf8.decode) .through(text.lines) .map(line => line) .intersperse("\n") .through(text.utf8.encode) .through(io.stdout) } val run = doThing(Path("t")).compile.drain }
但不清楚如何将GZIP解压整合到fs2的流处理中,请问该如何实现?
解决方案
只需要在读取文件字节流后,通过compression.gunzip()这个Pipe完成解压,再继续后续的文本解码、分行等操作即可。修改后的完整代码如下:
import cats.effect._ import fs2.io.file.{Files, Path} import fs2.{Stream, text, compression, io} object Main extends IOApp.Simple { def doThing(inPath: Path): Stream[IO, Unit] = { Files[IO] .readAll(inPath) // 新增:对GZIP压缩字节流进行解压 .through(compression.gunzip()) .through(text.utf8.decode) .through(text.lines) .map(line => line) .intersperse("\n") .through(text.utf8.encode) .through(io.stdout) } val run = doThing(Path("t.gz")).compile.drain }
关键说明
compression.gunzip()是fs2提供的流式解压Pipe,接收Stream[IO, Byte](压缩文件的字节流),输出解压后的Stream[IO, Byte],完全契合fs2的流式处理模型。- 整个流程保持了函数式的纯性,并且fs2会自动处理资源的打开和释放,无需手动管理输入流,避免了传统IO的资源泄漏问题。
- 后续的文本解码、分行、输出逻辑和处理普通文件完全一致,无需额外修改。
内容的提问来源于stack exchange,提问作者Jeremie
相关产品推荐
相关产品推荐

