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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 06:45:03