Fs2多管道间异常处理:单文件出错不中断整体且准确定位错误位置
解决方案
核心思路是通过阶段错误打标+单流统一错误捕获+完整性保障机制同时满足三个需求,不需要对单步流做多次编译,也不会出现多层错误处理逻辑被触发的问题。
第一步:实现阶段错误标记工具
不需要在每个步骤后加attempt/rethrow,只需要给每个处理管道包装一层错误适配逻辑,自动给异常附加阶段标识,后续捕获错误时可以直接定位出错位置:
import cats.MonadError import fs2.{Pipe, Stream} import java.nio.file.{Path, Files, StandardCopyOption} // 给处理管道添加阶段错误标记,出错时自动带上阶段信息 def withStage[F[_], I, O](stageName: String, pipe: Pipe[F, I, O])(implicit F: MonadError[F, Throwable]): Pipe[F, I, O] = input => pipe(input).adaptError { case ex => new RuntimeException(s"处理阶段失败:$stageName", ex) }
第二步:单文件处理流完整性保障
针对文件处理要求全成功/全失败的特性,采用临时文件中转机制:所有中间结果先写入临时文件,只有整个流处理完成无异常,才原子性移动到目标位置,中途出错自动清理临时文件,避免产生不完整的无效文件。单文件流只在最外层做一次错误捕获,不会触发多层处理逻辑:
// 单个文件的完整处理流 def processSingleFile[F[_]: MonadError[*[_], Throwable]](path: Path, tempDir: Path, targetDir: Path): Stream[F, Unit] = { val tempPath = tempDir.resolve(path.getFileName) val targetPath = targetDir.resolve(path.getFileName) // bracket 保障出错时自动清理临时文件 Stream.bracket(Stream.unit)(_ => Stream.eval(F.delay(Files.deleteIfExists(tempPath)))) *> downloadFile(path) // 给下载步骤单独加错误打标 .adaptError { case ex => new RuntimeException("处理阶段失败:文件下载", ex) } .through(withStage("病毒扫描", scanForViruses)) .through(withStage("内容加密", encrypt)) // 先写入临时文件 .through(Files.writeAll(tempPath)) // 全流程无错误则原子移动到最终目标位置 .evalMap(_ => F.delay(Files.move(tempPath, targetPath, StandardCopyOption.ATOMIC_MOVE))) // 统一错误捕获,仅触发一次,错误自带阶段信息 .handleErrorWith(ex => Stream.eval(log.error(s"文件${path.getFileName}处理失败", ex)) >> Stream.empty ) }
第三步:整体处理流组装
直接按原有逻辑平铺即可,单个文件出错不会影响其他文件的处理:
def processFiles[F[_]: MonadError[*[_], Throwable]](tempDir: Path, targetDir: Path): Stream[F, Unit] = getFilePaths.flatMap(path => processSingleFile(path, tempDir, targetDir))
小文件场景优化
如果处理的都是小文件,不想落盘生成临时文件,也可以用内存缓存替代,收集全量处理结果确认无异常后再向下游发送:
def processSingleFileInMemory[F[_]: MonadError[*[_], Throwable]](path: Path): Stream[F, Unit] = downloadFile(path) .adaptError { case ex => new RuntimeException("处理阶段失败:文件下载", ex) } .through(withStage("病毒扫描", scanForViruses)) .through(withStage("内容加密", encrypt)) // 收集全量字节确认无错误后再向下游发送 .compile.toVector .flatMap(bytes => Stream.emits(bytes)) .through(withStage("持久化存储", saveElsewhere)) .handleErrorWith(ex => Stream.eval(log.error(s"文件${path.getFileName}处理失败", ex)) >> Stream.empty )
方案优势
- 错误定位精准:所有异常都携带阶段标识,不需要多层handler即可直接判断出错步骤
- 错误处理逻辑仅触发一次:单文件流仅最外层有一个错误捕获点,不会出现多层handler依次执行的问题
- 无无效内容输出:通过临时文件或内存缓存保证只有完整处理成功的文件才会被保留
- 流仅编译一次:整个处理链是连续的流操作,没有多次编译的额外开销
内容的提问来源于stack exchange,提问作者Henry Parker
相关产品推荐
相关产品推荐

