Scala Spark:如何合并foreach中读取的多DataFrame为单个DataFrame
解决方案:合并不同结构CSV提取指定列后的DataFrames
当然有办法实现你的需求!你当前的代码是循环读取每个文件后单独处理,但没有把生成的DataFrames保存下来进行合并。只需要稍微调整代码逻辑,就能在内存中直接合并这些DataFrames,完全不需要写入文件再重新读取。
具体步骤:
- 保留原有的列定义与文件路径:
val cols = List("ID", "Subject") val filePaths = List("path to a.csv", "path to b.csv")
- 批量读取并筛选列,收集所有DataFrames:
不再用foreach(它只是执行操作但不返回结果),而是用map把每个文件对应的筛选后DataFrame收集到一个序列中:
val individualDFs = filePaths.map(path => { spark.read .option("header", "true") .option("delimiter", ",") .csv(path) .select(cols.head, cols.tail: _*) })
- 合并所有DataFrames:
使用reduce结合union操作,把序列中的所有DataFrames依次合并成一个:
val combinedDF = individualDFs.reduce(_ union _)
验证结果:
执行combinedDF.show()就能得到你想要的合并后结果:
+---+--------+ | ID| Subject| +---+--------+ | 1| English| | 2| IT| | 3| Science| | 4| IT| +---+--------+
执行combinedDF.count()会返回4,符合预期。
注意:如果你的Spark版本低于2.0,需要用
unionAll代替union,不过Spark 2.0及以上版本中union已经是unionAll的别名,功能完全一致。
内容的提问来源于stack exchange,提问作者CRV
相关产品推荐
相关产品推荐

