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

Akka Streams故障时跳过Flow且不丢失数据的实现方法

Akka Streams处理失败时保留原消息的解决方案

你的核心需求是在Flow处理元素失败时,不丢失该元素,直接将原元素传递到下游。当前代码的问题在于:

  • map中抛出异常后,监督策略的Supervision.restart会重启Flow,但当前失败的元素(5)会被直接丢弃
  • recover并未生效,因为监督策略阻止了流的终止,recover只在流崩溃时触发

直接解决方案:在处理逻辑内部捕获异常

最简洁的方式是在map的处理逻辑里捕获异常,当处理失败时直接返回原元素,避免抛出异常导致元素丢失:

Source(1 to 10)
  .via(Flow[Int].map { x =>
    try {
      if (x == 5) throw new Exception("boom!")
      x
    } catch {
      case _: Exception => x // 处理失败时返回原元素
    }
  })
  .runWith(Sink.foreach(println))

运行这段代码会输出你期望的结果:1 2 3 4 5 6 7 8 9 10

为什么原方案不生效?

原代码中,当x=5抛出异常时,监督策略的restart会重启当前Flow的处理逻辑,但已经进入处理流程的元素5会被丢弃,不会传递到下游。而recover操作只有在整个流发生终止级别的错误时才会触发,这里因为监督策略的存在,流并未终止,所以recover完全没起作用。

替代方案:用Either包装处理结果

如果需要区分正常处理的元素和失败的元素(后续可以做额外处理),可以用Either包装结果,再展开传递原元素:

Source(1 to 10)
  .via(Flow[Int].map { x =>
    try {
      Right(if (x == 5) throw new Exception("boom!") else x)
    } catch {
      case e: Exception => Left((x, e))
    }
  })
  .map { // 展开Either,失败时取原元素
    case Right(value) => value
    case Left((original, _)) => original
  }
  .runWith(Sink.foreach(println))

这种方式既保留了原元素,也可以在中间步骤记录错误信息,适合需要监控失败情况的场景。

内容的提问来源于stack exchange,提问作者David Santiago Gantiva Castro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 15:06:14