Akka Streams如何按字节大小对源中的文件列表进行分块?
我明白你的痛点——处理GB级别的大文件时,用Sink.fold把所有数据攒在内存里绝对是个灾难,而且你试过的grouped(N)是按元素数量分块,完全没法解决字节大小统一的需求。别担心,Akka Streams有专门的方式来实现按字节大小分块的流处理,全程不会把整个文件加载到内存里。
核心思路
咱们要把原来的“先攒所有数据再写入”改成“边读边切边写”:
- 读取多个文件的字节流并合并成一个连续流
- 把这个连续流切成固定字节大小的块(比如10MB/块)
- 逐个块异步写入输出文件(或分块文件)
实现步骤
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

