Akka Stream如何按顺序串联notify与write操作且不丢失OutputRecord
解决方案
你遇到的报错本质是notify Flow的输出类型为PushResult,和下游writeOutput要求的输入类型不匹配,且执行通知逻辑后丢失了后续写入需要的OutputRecord对象。由于要求先执行通知成功再写入的顺序逻辑,只需要在通知执行完成后将原始OutputRecord透传给下游即可实现需求,无需并行处理。
具体实现代码
import akka.stream.scaladsl.Source import akka.util.ByteString readAsCSV.flatMap { recordSource => recordSource // 第一步:将读取到的CSV键值对转为InputRecord,补全你之前遗漏的转换逻辑 .map { csvMap => InputRecord( recordId = csvMap("recordId").utf8String, name = csvMap("name").utf8String, salary = csvMap("salary").utf8String.toLong ) } .map(process) // 转为OutputRecord // 核心逻辑:执行通知,成功后透传原始OutputRecord到下游 .flatMapConcat { outputRecord => // 单元素流走通知逻辑,等通知返回结果后再往下传递原始OutputRecord Source.single(outputRecord).via(notify).map(_ => outputRecord) } // 将OutputRecord序列化为CSV格式的ByteString,适配写入接口的输入要求 .map { outputRecord => ByteString(s"${outputRecord.recordId},${outputRecord.name},${outputRecord.designation}\n") } .to(writeOutput) .run() }
逻辑说明
flatMapConcat会等待每一条记录的notify执行完成拿到PushResult后,才会将原始OutputRecord传递到下游,完全满足顺序执行要求,不会出现并行处理的问题- 如果需要用到
PushResult里的附加信息,只要把.map(_ => outputRecord)改为.map(pushResult => (outputRecord, pushResult)),就能将两个值一起带到下游做自定义处理 - 如果你可以修改
notify方法的定义,也可以直接把它的返回值改为Flow[OutputRecord, OutputRecord, NotUsed],内部处理完通知逻辑后返回原始输入的OutputRecord,流的拼接会更简洁。
内容的提问来源于stack exchange,提问作者Aiden
相关产品推荐
相关产品推荐

