You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何处理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
  }
}

关键逻辑说明

两种方式本质都是隔离单个元素的异常影响:

  1. 捕获id=3处理时的异常,不会让异常扩散到整个流
  2. 过滤掉失败的元素,让流继续处理后续的1、2、3循环中的有效元素
  3. 流会持续运行,无限输出"name Olivia"和"name Liam"

内容的提问来源于stack exchange,提问作者Randhir Kumar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 20:58:11