如何将Akka Source直接拆分为批量子Source以优化内存?
Akka Stream 无界 Source 批量拆分优化方案
问题场景
我有一个文件解析器,定义如下:
def parser(file: File): Source[Record, NotUsed] = ...
还有一个处理Record的函数:
def send(source: Source[Record, _]): Future[Result]
由于文件体积较大,且send函数对单次可流式处理的记录数量有限制,我需要将文件对应的Source拆分为批量。目前通过内存缓冲区实现的代码如下:
parser(file) .grouped(batchSize) .mapAsync(2) { (batch: Seq[Record]) => send(Source(batch)) }
该实现可正常工作,但如果能直接将无界Source拆分为子Source,无需使用中间的Seq[Record]作为载体,就能进一步优化性能、降低内存占用。是否存在这样的实现方式?比如类似:
parser(file) .groupedToSources(batchSize) .mapAsync(2) { (batch: Source[Record, NotUsed]) => send(batch) }
本质上我需要一个等效于以下逻辑的操作符,但不需要依赖内存中的Vector[Record]做中间存储:
Flow[Record].grouped(batchSize).map(Source _)
解决方案:自定义 groupedToSources 操作符
Akka Stream 原生未提供该操作,但可以通过自定义GraphStage实现,核心思路是逐个接收上游元素,累计到指定数量后生成子Source发送,全程无需将整个批次加载到内存集合中。
自定义实现代码
import akka.stream._ import akka.stream.stage._ import scala.concurrent.Promise import akka.NotUsed class GroupedToSourceStage[T](batchSize: Int) extends GraphStage[FlowShape[T, Source[T, NotUsed]]] { require(batchSize > 0, "batchSize must be positive") val in: Inlet[T] = Inlet[T]("GroupedToSource.in") val out: Outlet[Source[T, NotUsed]] = Outlet[Source[T, NotUsed]]("GroupedToSource.out") override val shape: FlowShape[T, Source[T, NotUsed]] = FlowShape(in, out) override def createLogic(inheritedAttributes: Attributes): GraphStageLogic = new GraphStageLogic(shape) { private var currentBatchCount = 0 private var currentSourceBuffer: List[T] = Nil setHandler(in, new InHandler { override def onPush(): Unit = { val elem = grab(in) currentBatchCount += 1 currentSourceBuffer = elem :: currentSourceBuffer if (currentBatchCount == batchSize) { // 生成子Source并推送下游 push(out, Source(currentSourceBuffer.reverse)) // 重置批次状态 currentBatchCount = 0 currentSourceBuffer = Nil } else { pull(in) } } override def onUpstreamFinish(): Unit = { // 处理上游结束后剩余的不足一批次的元素 if (currentBatchCount > 0) { emit(out, Source(currentSourceBuffer.reverse)) } completeStage() } }) setHandler(out, new OutHandler { override def onPull(): Unit = { if (currentBatchCount == 0) { pull(in) } else { // 推送当前已累计的批次 push(out, Source(currentSourceBuffer.reverse)) currentBatchCount = 0 currentSourceBuffer = Nil } } }) } } // 扩展方法简化调用 object GroupedToSource { def apply[T](batchSize: Int): Flow[T, Source[T, NotUsed], NotUsed] = { Flow.fromGraph(new GroupedToSourceStage[T](batchSize)) } }
使用方式
直接用自定义操作符替代原有的grouped+map(Source _)组合即可:
parser(file) .via(GroupedToSource(batchSize)) .mapAsync(2) { (batchSource: Source[Record, NotUsed]) => send(batchSource) }
核心原理
- 自定义
GraphStage逐个接收上游元素,累计到指定批次数量时,将缓存的元素反转后封装为子Source推送下游 - 上游流结束时,自动处理剩余的不足一批次的元素,避免数据丢失
- 全程仅缓存当前批次的元素,且无额外内存集合拷贝,有效降低大文件场景下的内存消耗
内容的提问来源于stack exchange,提问作者ArunasR
相关产品推荐
相关产品推荐

