FS2错误处理实践有哪些?现有代码求更优惯用实现方案
FS2 错误处理的惯用写法优化
针对你这段温度转换的代码,我们可以利用FS2和Cats Effect的原生错误处理机制替代scala.util.Try,让代码更贴合函数式流处理的惯用风格,同时保持逻辑清晰。
核心优化思路
- 抛弃
Try,改用Cats EffectIO的异常处理能力——FS2 Stream本身基于IO,天然支持错误传播与捕获,无需手动用Try包裹。 - 使用FS2的
attempt算子将流元素转换为Either[Throwable, T],明确区分成功与失败的元素,便于分别处理。 - 避免在
map里混合成功结果和错误信息,改用更清晰的流分支逻辑。
优化后的代码(保留原需求:错误转字符串写入文件)
import cats.effect.{IO, IOApp} import fs2.{Stream, text} import fs2.io.file.{Files, Path} object Converter extends IOApp.Simple { private val converter: Stream[IO, Unit] = { def fahrenheitToCelsius(f: Double): Double = (f - 32.0) * (5.0 / 9.0) // 将转换逻辑包装为IO,自动捕获同步异常 def convertLine(line: String): IO[String] = IO.delay { fahrenheitToCelsius(line.toDouble).toString } Files[IO].readUtf8Lines(Path("testdata/fahrenheit.txt")) .filter(s => s.trim.nonEmpty && !s.startsWith("//")) // 执行转换并捕获异常,转为Either类型 .evalMap(line => convertLine(line).attempt) // 统一处理成功/失败结果,转为字符串输出 .map { case Right(celsiusStr) => celsiusStr case Left(cause) => s"Failed to convert $line: ${cause.getMessage}" } .intersperse("\n") .through(text.utf8.encode) .through(Files[IO].writeAll(Path("testdata/celsius.txt"))) } def run: IO[Unit] = converter.compile.drain }
更灵活的错误处理方案(分离成功输出与错误日志)
如果想把错误单独记录日志,只将成功结果写入文件,可以这样调整:
import cats.effect.{IO, IOApp} import fs2.{Stream, text} import fs2.io.file.{Files, Path} object Converter extends IOApp.Simple { private val converter: Stream[IO, Unit] = { def fahrenheitToCelsius(f: Double): Double = (f - 32.0) * (5.0 / 9.0) def convertLine(line: String): IO[Double] = IO.delay { fahrenheitToCelsius(line.toDouble) } val lines = Files[IO].readUtf8Lines(Path("testdata/fahrenheit.txt")) .filter(s => s.trim.nonEmpty && !s.startsWith("//")) // 拆分成功转换与错误处理 val successfulConversions = lines .evalMap(line => convertLine(line).attempt) .flatMap { case Right(celsius) => Stream.emit(celsius.toString) case Left(cause) => // 记录错误日志,不中断流的执行 Stream.eval(IO.println(s"Conversion error: ${cause.getMessage}")) >> Stream.empty } successfulConversions .intersperse("\n") .through(text.utf8.encode) .through(Files[IO].writeAll(Path("testdata/celsius.txt"))) } def run: IO[Unit] = converter.compile.drain }
关键知识点说明
evalMap+IO.attempt:evalMap用于执行返回IO的操作,IO.attempt会将IO的成功结果包装为Right,异常包装为Left,避免单个元素失败导致整个流中断。- 流分支处理:通过
flatMap可以根据Either结果分支处理,比如错误时输出日志并跳过该元素,成功时保留结果。 - 避免手动try/catch:IO的
delay方法会自动捕获同步异常,无需手动编写try/catch块,代码更简洁。
内容的提问来源于stack exchange,提问作者Eric Kolotyluk
相关产品推荐
相关产品推荐

