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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 14:27:23