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

Akka Streams如何按字节大小对源中的文件列表进行分块?

按字节大小分块处理大文件流(Akka Streams)

我明白你的痛点——处理GB级别的大文件时,用Sink.fold把所有数据攒在内存里绝对是个灾难,而且你试过的grouped(N)是按元素数量分块,完全没法解决字节大小统一的需求。别担心,Akka Streams有专门的方式来实现按字节大小分块的流处理,全程不会把整个文件加载到内存里。

核心思路

咱们要把原来的“先攒所有数据再写入”改成“边读边切边写”:

  1. 读取多个文件的字节流并合并成一个连续流
  2. 把这个连续流切成固定字节大小的块(比如10MB/块)
  3. 逐个块异步写入输出文件(或分块文件)

实现步骤

1. 实现固定字节大小的分块Flow

首先需要一个自定义Flow,它能把任意大小的ByteString流,切成我们指定字节大小的块。我们用statefulMapConcat来维护一个累加器,确保内存里永远只保留未凑够块大小的剩余字节:

import akka.stream.scaladsl.Flow
import akka.util.ByteString
import akka.NotUsed

// 定义分块大小,比如10MB
val chunkSize = 1024 * 1024 * 10 

val fixedSizeChunkFlow: Flow[ByteString, ByteString, NotUsed] =
  Flow[ByteString].statefulMapConcat { () =>
    // 每个流实例维护自己的累加器,线程安全
    var accumulator = ByteString.empty
    
    { incomingBs =>
      // 把新收到的字节加到累加器里
      accumulator ++= incomingBs
      // 把累加器切成指定大小的块
      val chunks = accumulator.grouped(chunkSize).toList
      // 留下最后一块(可能不够大小),下次继续累加
      accumulator = chunks.lastOption.getOrElse(ByteString.empty)
      // 输出所有凑够大小的块
      chunks.init
    }
  }
  // 流结束时,把最后剩余的字节(如果有的话)输出
  .concatMat(Source.single(ByteString.empty).filter(_.nonEmpty))(Keep.left)

2. 重构流处理流程

接下来替换掉原来的Sink.fold,用分块Flow+异步写入Sink来处理:

场景1:所有分块写入同一个文件(追加模式)

用Akka的FileIO.toPath做异步写入,比Java原生的Files.write更适合流处理:

import akka.stream.scaladsl.{Source, FileIO}
import java.nio.file.{Paths, StandardOpenOption}

val files = List("a.txt", "b.txt", "c.txt")
val outputPath = Paths.get("an-output-file.txt")

// 构建写入Sink:创建文件(如果不存在),追加模式
val appendSink = FileIO.toPath(
  outputPath,
  Set(StandardOpenOption.CREATE, StandardOpenOption.APPEND)
)

// 组装整个流
Source(files)
  // 读取每个文件的字节流,合并成连续流
  .flatMapConcat(file => FileIO.fromPath(Paths.get(file)))
  // 切成固定大小的块
  .via(fixedSizeChunkFlow)
  // 逐个块写入文件
  .runWith(appendSink)

场景2:每个分块写入单独的文件

如果需要把每个块存成独立文件(比如output-0.txt、output-1.txt),可以用zipWithIndex给块编号,再异步写入:

Source(files)
  .flatMapConcat(file => FileIO.fromPath(Paths.get(file)))
  .via(fixedSizeChunkFlow)
  // 给每个块加上索引编号
  .zipWithIndex
  // 异步写入(控制并行度,避免IO压力过大)
  .mapAsync(parallelism = 4) { case (chunk, index) =>
    val chunkPath = Paths.get(s"an-output-file-$index.txt")
    FileIO.write(chunkPath, chunk).run()
  }
  .runWith(Sink.ignore)

关键细节解释

  • 为什么不用grouped(N)?
    grouped(N)是按元素数量分块,而不是字节大小。比如如果某个文件的ByteString是2GB,grouped(1)直接就会把2GB加载到内存里,完全达不到分块目的。我们的自定义Flow是按字节总大小切分,不管单个ByteString的大小,确保每个输出块的字节数不超过设定值。

  • 内存占用控制?
    整个流程中,内存里最多只会保留一个未凑够块大小的剩余字节(最多chunkSize-1字节),加上当前正在写入的块,内存占用完全可控,不会出现OOM。

内容的提问来源于stack exchange,提问作者Ivan Peng

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:09:52