Fs2流转换为Zip压缩流时无法终止的问题排查
问题分析与解决方案
首先,你的核心问题出在enqueueChunkSync方法里的阻塞式队列操作——它用了SyncVar和unsafeRunSync,完全违背了cats-effect的异步非阻塞模型,尤其是在处理最后一个块时,这种阻塞逻辑很容易导致死锁。
为什么最后一个块会特殊?
前面的块都是在正常写入zip条目的流程中产生的,此时write流(处理zip条目并写入数据的流)还在活跃运行,和read流处于concurrently的协同状态,队列的读写是顺畅的。
但最后一个块是在ZipOutputStream.close()时触发的:这属于bracket的释放阶段,此时write流已经完成了所有条目的处理,正在执行收尾操作。这时候用阻塞的方式往队列里塞数据,很可能因为线程上下文的冲突(比如当前线程是阻塞线程池的线程,而队列的enqueue操作需要等待上下文切换),导致整个流程卡住,无法终止。
你提到“移除阻塞代码后功能看似正常”,其实是因为去掉阻塞后,队列操作回到了异步模型里,符合fs2和cats-effect的设计原则,并不会导致“字节一次性刷新”——fs2本身会帮你处理流的分块和背压。
修复方案:用异步方式实现OutputStream
我们需要把自定义的OutputStream改成纯异步实现,完全抛弃阻塞式的SyncVar和unsafeRunSync,改用cats-effect的Async能力处理队列写入。
修改后的核心代码如下:
private def zipP1[F[_]](implicit F: Async[F], blockingEc: ExecutionContext, contextShift: ContextShift[F]): Pipe[F, (String, Stream[F, Byte]), Byte] = entries => { Stream.eval(Queue.unbounded[F, Option[Chunk[Byte]]]).flatMap { q => Stream.suspend { val os = new java.io.OutputStream { private var bufferedChunk: Chunk[Byte] = Chunk.empty private def enqueueChunk(a: Option[Chunk[Byte]]): Unit = { // 用Async的async方法把异步操作桥接到OutputStream的同步API里 F.async_[Unit] { cb => q.enqueue1(a).attempt.flatMap { case Right(_) => F.delay(cb(Right(()))) case Left(e) => F.delay(cb(Left(e))) }.unsafeRunAsyncAndForget() }.unsafeRunSync() } @scala.annotation.tailrec private def addChunk(c: Chunk[Byte]): Unit = { val free = 1024 - bufferedChunk.size if (c.size > free) { enqueueChunk(Some(Chunk.vector(bufferedChunk.toVector ++ c.take(free).toVector))) bufferedChunk = Chunk.empty addChunk(c.drop(free)) } else { bufferedChunk = Chunk.vector(bufferedChunk.toVector ++ c.toVector) } } override def close(): Unit = { if (bufferedChunk.nonEmpty) { enqueueChunk(Some(bufferedChunk)) bufferedChunk = Chunk.empty } enqueueChunk(None) } override def write(bytes: Array[Byte]): Unit = addChunk(Chunk.bytes(bytes)) override def write(bytes: Array[Byte], off: Int, len: Int): Unit = addChunk(Chunk.bytes(bytes, off, len)) override def write(b: Int): Unit = addChunk(Chunk.singleton(b.toByte)) } val write: Stream[F, Unit] = Stream .bracket(F.delay(new ZipOutputStream(os)))((zos: ZipOutputStream) => F.delay(zos.close())) .flatMap((zos: ZipOutputStream) => entries.through(writeEntry(zos))) val read = q.dequeue .unNoneTerminate .flatMap(Stream.chunk(_)) read.concurrently(write) } } }
关键改进点:
- 用
F.async_替代SyncVar,把队列的异步enqueue操作安全桥接到OutputStream的同步API中,避免阻塞线程。 - 去掉所有
unsafeRunSync的阻塞调用,改用unsafeRunAsyncAndForget()触发异步操作,不会卡住当前线程。 - 保留原有缓冲区逻辑,确保分块写入行为和之前一致,fs2的背压机制会自动控制数据流动,不会出现“一次性刷新字节”的问题。
这样修改后,最后一个块的写入会和其他块一样在异步模型中顺畅完成,流也能正常终止。
内容的提问来源于stack exchange,提问作者Leonti
相关产品推荐
相关产品推荐

