Cats Effect IO flatMap内用attempt/redeem无法捕获异常解决方案
Cats Effect 中 flatMap 内部调用 attempt/redeem 无法捕获异常的问题解决
问题背景
在由多个小型IO组合而成的程序中传递事件数据时,需要基于事件执行可能抛出异常的计算逻辑,后续上报执行结果时需包含原始事件信息。
实现过程中观察到不符合预期的行为:
- 初始实现思路是先执行计算,再通过
attempt/redeem方法将抛出的异常转换为值处理 - 在
flatMap调用内部使用attempt(或基于attempt实现的redeem)时,异常不会被捕获,会直接导致整个IOApp崩溃 - 如果将
attempt/redeem放在应用最顶层调用,异常可以正常被捕获并转换为值,但此时成功、失败的处理逻辑无法获取原始事件对象 - 并行执行场景存在相同问题:例如调用
.parTupled(program1.attempt, program2.attempt)时,任意程序抛出异常都会导致应用崩溃 - 已知可通过Reader、Kleisli等方式实现数据传递,但对于当前需求会引入不必要的额外开销
问题复现代码
import cats.effect.{IO, IOApp} object ParallelExecutionWAttempt extends IOApp.Simple { def run: IO[Unit] = mainWInnerRedeem /** redeem放在flatMap内部的实现 * * 预期:所有抛出的异常被捕获为值后处理 * 实际:异常未被转换为值,直接导致整个应用崩溃 * */ def mainWInnerRedeem: IO[Unit] = getEventFromSource .flatMap{ event => getEventHandler(event).redeem(ex => onFailure(ex, event), _ => onSuccess(event)) } /** redeem放在调用链最外层的实现 * * 表现符合预期:异常可以被正常捕获 * 问题:成功/失败处理逻辑无法获取原始event对象 */ def mainWOuterRedeem: IO[Unit] = getEventFromSource.flatMap(getEventHandler) .redeem( ex => IO.println(s"程序执行失败,异常信息:$ex"), _ => IO.println("程序执行成功!") ) /** 演示用事件定义 */ trait Event case class Event1(a: Int) extends Event case class Event2(b: String) extends Event /** 事件源:返回包装在IO中的事件对象 */ def getEventFromSource: IO[Event] = IO{Event1(1)} /** 获取事件对应的处理器 */ def getEventHandler(event: Event): IO[Unit] = blowsUp(event) /** 测试用处理器:一个直接抛出异常,一个正常执行 */ def blowsUp(event: Event): IO[Unit] = throw new RuntimeException("执行抛出异常!") def successfulFunc(event: Event): IO[Unit] = IO{println("执行正常")} /** 结果处理函数:需要传入原始事件作为参数 */ def onSuccess(event: Event): IO[Unit] = IO.println(s"执行成功,事件内容:$event") def onFailure(throwable: Throwable, event: Event): IO[Unit] = IO.println(s"执行失败,异常:$throwable,事件内容:$event") }
问题根因
核心原因有两点:
- Scala采用严格求值策略,
attempt/redeem是IO实例的成员方法,只有成功构造出IO实例之后,调用它的redeem方法才能捕获这个IO实例运行时抛出的异常。如果在构造IO实例的过程中(也就是调用返回IO的方法时)就直接抛出异常,此时还没执行到redeem调用,异常会直接逃到IO的错误捕获机制之外。 - 示例代码中的
blowsUp方法属于错误写法:方法体直接抛出异常,异常发生在IO构造阶段而非IO运行阶段,内部调用的redeem根本没有执行机会。 - 外层redeem能捕获异常的原因是:IO原生的
flatMap、map等组合子本身被IO运行时的错误机制包裹,执行传入的函数时如果抛出异常,会被自动转换为IO的失败值向上传递,最终被外层的redeem捕获。
解决方案
1. 规范IO方法的异常抛出方式(推荐)
所有返回IO的方法,不要在方法体中直接抛出异常,要把异常逻辑放到IO的求值逻辑中,有两种标准写法:
// 写法1:用IO.raiseError直接返回包含异常的IO实例 def blowsUp(event: Event): IO[Unit] = IO.raiseError(new RuntimeException("执行抛出异常!")) // 写法2:把可能抛异常的逻辑放到IO代码块中,IO.apply会自动捕获块内抛出的异常 def blowsUp(event: Event): IO[Unit] = IO { throw new RuntimeException("执行抛出异常!") }
修改后无论redeem放在flatMap内部还是外层,都能正常捕获异常,且内部redeem可以直接拿到闭包中的原始event对象,完全符合需求。
2. 无法修改第三方方法时用IO.defer包裹
如果调用的是第三方返回IO的方法,无法修改其实现,不确定是否会在构造阶段抛异常,可以用IO.defer包裹方法调用,它会把构造IO时抛出的异常自动捕获,转换为IO内部的错误值:
def mainWInnerRedeem: IO[Unit] = getEventFromSource .flatMap{ event => // IO.defer会捕获getEventHandler构造IO时抛出的异常 IO.defer(getEventHandler(event)) .redeem(ex => onFailure(ex, event), _ => onSuccess(event)) }
并行场景修复
之前提到的parTupled场景问题根因完全一致:传入的program参数在求值时就抛出了异常,没有被包装到IO运行时逻辑中。修复方式相同,要么修改对应IO方法的异常抛出逻辑,要么用IO.defer包裹每个program的构造过程:
// 错误写法:如果program1/program2在调用时直接抛异常,会直接崩溃 val badCase = (program1.attempt, program2.attempt).parTupled // 正确写法:IO.defer捕获构造阶段的异常 val goodCase = (IO.defer(program1).attempt, IO.defer(program2).attempt).parTupled
内容的提问来源于stack exchange,提问作者six_minute_abs
相关产品推荐
相关产品推荐

