Akka Streams:如何用GraphDSL构建Source of Sources?
Akka Stream 多分流后生成ZIP压缩包的解决方案
问题背景
从数据库获取行数据得到Akka Source,需按分区函数分流到不同流,通过GraphDSL实现扇出:
- 先创建
UniformFanOutShape类型的Partition,连接原始Source入口,Partition出口悬空 - 每个分流流用Alpakka CSV模块的
CsvFormatting.format()做格式化,连接到Partition出口,Flow出口仍悬空 - 需将每个CSV流添加文件名并打包成ZIP,Alpakka Archive的
Archive.zip()要求输入Source[(ArchiveMetadata, Source[ByteString, Any]), _],但当前图有多个悬空出口,无法直接转换为符合要求的Sources集合;尝试自定义多Graph生成方法,但Akka内部API为包私有无法直接使用,且单个SourceShape会因未连接所有Partition出口触发图验证报错。
可行解决方案
方案1:封装多出口Graph,逐个提取Source
将分流+CSV格式化的逻辑封装为无入口、多出口的AmorphousShape Graph,再通过Source.fromGraph逐个提取每个出口对应的Source,最后包装成Archive所需的输入格式:
import akka.stream._ import akka.stream.alpakka.csv.scaladsl.CsvFormatting import akka.stream.alpakka.file.scaladsl.Archive import akka.stream.scaladsl._ import akka.util.ByteString import java.nio.file.Paths // 假设Row是数据库行数据类型 case class Row(/* 行数据字段 */) val originalSource: Source[Row, _] = ??? // 数据库获取的Source val partitionCount = 3 // 分流数量 // 构建包含分流和CSV格式化的多出口Graph val multiOutputGraph = GraphDSL.create() { implicit builder => import GraphDSL.Implicits._ // 创建Partition扇出节点 val partition = builder.add(Partition[Row](partitionCount, row => { // 自定义分区逻辑:返回0到partitionCount-1的索引 row./* 字段 */ % partitionCount })) // 创建对应数量的CSV格式化Flow val csvFlows = (0 until partitionCount).map(_ => builder.add(CsvFormatting.format(/* 传入CSV格式化参数,如分隔符、表头等 */)) ) // 连接Partition到每个CSV Flow (0 until partitionCount).foreach(i => partition.out(i) ~> csvFlows(i)) // 返回无入口、多出口的AmorphousShape AmorphousShape(inputPorts = Seq.empty, outputPorts = csvFlows.map(_.out)) } // 从多出口Graph中提取每个CSV流的Source val csvSources: Seq[Source[ByteString, _]] = (0 until partitionCount).map { idx => Source.fromGraph(multiOutputGraph).out(idx) } // 包装成Archive要求的输入格式:(ArchiveMetadata, Source[ByteString, Any]) val archiveInput = Source( csvSources.zipWithIndex.map { case (src, idx) => ArchiveMetadata(s"partition_$idx.csv") -> src } ) // 生成ZIP压缩流并写入文件 val zipSource = archiveInput.via(Archive.zip()) zipSource.runWith(FileIO.toPath(Paths.get("partitioned_data.zip")))
方案2:用groupBy替代GraphDSL(分区逻辑适用时)
如果分区逻辑可以通过分组键实现,直接使用groupBy分流,无需GraphDSL,更简洁:
val zipSource = originalSource // 按分区键分组,分流为多个子流 .groupBy(partitionCount, row => row./* 分区键字段 */) // 每个子流做CSV格式化 .via(CsvFormatting.format(/* 格式化参数 */)) // 将每个子流的ByteString收集为单个Source .fold(Source.empty[ByteString]) { (acc, bs) => acc ++ Source.single(bs) } // 合并所有子流的结果 .mergeSubstreams // 为每个Source添加文件名元数据 .zipWithIndex .map { case (src, idx) => ArchiveMetadata(s"group_$idx.csv") -> src } // 生成ZIP流 .via(Archive.zip()) // 写入文件 zipSource.runWith(FileIO.toPath(Paths.get("grouped_data.zip")))
关于包私有API的说明
Akka Stream将GraphDSL.Builder和createGraph设为包私有,是为了保证流图的完整性约束,避免不完整的流定义导致运行时错误。强行绕过会破坏Akka的图验证机制,引发未连接端口等报错,不建议尝试。
内容的提问来源于stack exchange,提问作者silverberry
相关产品推荐
相关产品推荐

