如何处理Scala FS2流转换异常以保障流不间断运行?
FS2流异常处理:保证流不间断运行
现有如下Scala FS2代码,当流处理id=3时,由于studentData中无对应数据会抛出异常,导致整个流终止。需求是让流持续输出"name Olivia"和"name Liam",不受异常影响,保持流无限运行。
import cats.effect.{IO, IOApp} import fs2.Pipe import fs2.Stream object Test extends IOApp.Simple { final case class Student(id: Int, name: String) private val studentData: Map[Int, Student] = Map(1 -> Student(1, "Olivia"), 2 -> Student(2, "Liam")) override def run: IO[Unit] = { def firstTransformation: Pipe[IO, Int, Student] = { stream => stream.evalMap(id => IO(studentData(id))) // 此处会因id=3抛出异常 } def secondTransformation: Pipe[IO, Student, String] = { stream => stream.evalMap(student => IO(student.name)) } val sourceStream: Stream[IO, Int] = Stream(1, 2, 3).repeat sourceStream .through(firstTransformation) .through(secondTransformation) .map(name => println(s"name $name")) .compile .drain } }
解决方案
核心思路是捕获单个元素处理时的异常,过滤掉失败的元素,让流继续处理后续元素。提供两种实现方式:
方式一:用attempt包装IO结果并过滤
import cats.effect.{IO, IOApp} import fs2.Pipe import fs2.Stream object Test extends IOApp.Simple { final case class Student(id: Int, name: String) private val studentData: Map[Int, Student] = Map(1 -> Student(1, "Olivia"), 2 -> Student(2, "Liam")) override def run: IO[Unit] = { def firstTransformation: Pipe[IO, Int, Student] = { stream => // 将IO操作结果包装为Either,捕获异常 stream.evalMap(id => IO(studentData(id)).attempt) // 只保留成功获取的Student,过滤异常对应的Left值 .collect { case Right(student) => student } } def secondTransformation: Pipe[IO, Student, String] = { stream => stream.evalMap(student => IO(student.name)) } val sourceStream: Stream[IO, Int] = Stream(1, 2, 3).repeat sourceStream .through(firstTransformation) .through(secondTransformation) .map(name => println(s"name $name")) .compile .drain } }
方式二:在IO层直接处理异常并过滤
如果不想引入Either,可以用evalMapFilter直接返回可选结果:
import cats.effect.{IO, IOApp} import fs2.Pipe import fs2.Stream object Test extends IOApp.Simple { final case class Student(id: Int, name: String) private val studentData: Map[Int, Student] = Map(1 -> Student(1, "Olivia"), 2 -> Student(2, "Liam")) override def run: IO[Unit] = { def firstTransformation: Pipe[IO, Int, Student] = { stream => // 异常时返回None自动被过滤,成功时返回Some(Student) stream.evalMapFilter(id => IO(studentData(id)).handleErrorWith(_ => IO.pure(None)) ) } def secondTransformation: Pipe[IO, Student, String] = { stream => stream.evalMap(student => IO(student.name)) } val sourceStream: Stream[IO, Int] = Stream(1, 2, 3).repeat sourceStream .through(firstTransformation) .through(secondTransformation) .map(name => println(s"name $name")) .compile .drain } }
关键逻辑说明
两种方式本质都是隔离单个元素的异常影响:
- 捕获id=3处理时的异常,不会让异常扩散到整个流
- 过滤掉失败的元素,让流继续处理后续的1、2、3循环中的有效元素
- 流会持续运行,无限输出"name Olivia"和"name Liam"
内容的提问来源于stack exchange,提问作者Randhir Kumar
相关产品推荐
相关产品推荐

