实现Supervision.Resume的Akka Stream测试流技术咨询
我最近刚搞定了一个Akka Stream的流处理逻辑,专门用来处理JSON消息的解析与校验,还做了自定义的错误处理策略,分享给需要的朋友参考:
功能说明
这个流主要完成这几件核心事:
- 解析输入的JSON消息字符串
- 校验消息中是否存在
destination_region这个必填键 - 将原始消息和提取到的
destination_region值封装成样例类,传递给下游处理阶段 - 遇到解析失败或者键缺失的错误时,先记录异常日志,然后让流继续运行(不会因为单个错误中断整个流的处理)
最简实现代码
package com.example.stages import com.example.helpers.EitherHelper import akka.stream.scaladsl.Flow import akka.stream.{ActorAttributes, Supervision} import org.json4s._ import org.json4s.native.JsonMethods._ // 定义传递给下游的封装类,打包原始消息和目标区域值 case class ProcessedMessage(originalMsg: String, destinationRegion: String) object JsonProcessingStage { // 自定义错误决策器:捕获指定异常并恢复流处理 private val errorDecider: Supervision.Decider = { case e: JsonParsingException => println(s"JSON解析失败,错误详情: ${e.getMessage}") Supervision.Resume case e: NoSuchElementException => println(s"必填键缺失: destination_region不存在,错误详情: ${e.getMessage}") Supervision.Resume case _ => Supervision.Stop // 其他未定义的错误默认终止流 } val jsonProcessingFlow: Flow[String, ProcessedMessage, _] = Flow[String] .map(parse) // 解析输入的JSON字符串 .map { json => // 用EitherHelper安全提取目标键,提取失败则抛出异常交给decider处理 val region = EitherHelper.extract[String](json, "destination_region") .getOrElse(throw new NoSuchElementException("destination_region not found in JSON")) ProcessedMessage(compact(json), region) } .withAttributes(ActorAttributes.supervisionStrategy(errorDecider)) }
关键细节解析
- 自定义Supervision.Decider:专门针对JSON解析异常和键缺失异常做处理,打印错误日志后返回
Supervision.Resume,让流跳过当前错误消息、继续处理后续数据,避免单个坏消息导致整个流中断 - 样例类封装:
ProcessedMessage把原始JSON内容和提取到的区域值打包,下游阶段不需要重复解析JSON,直接使用封装好的数据即可 - EitherHelper辅助:这里用到的
EitherHelper.extract是自定义的工具方法,用来安全地从JSON对象中提取指定键的值,简化了Either类型的处理逻辑
内容的提问来源于stack exchange,提问作者mabe02
相关产品推荐
相关产品推荐

