Python背景用户如何在Scala中实现多个DataFrame的Union合并
Scala 多DataFrame批量Union解决方案
报错原因
你遇到的报错由两个问题导致:
- Scala中
val定义的是不可变变量,且你在if/else代码块内部定义的union_df仅在当前块的作用域生效,块外无法访问 - else分支中你尝试调用还未完成定义的
union_df,触发了递归引用报错
最优实现方案(推荐)
直接使用Scala集合的reduce方法实现批量union,不需要手动处理索引判断,代码更简洁高效,和Python实现的效果完全一致:
import org.apache.spark.sql.DataFrame // 先将所有需要合并的DataFrame存入List,对应Python侧的list_of_dfs val listOfDfs: List[DataFrame] = List( spark.createDataFrame(Seq(("A", "C"), ("B", "E"))).toDF("dummy1", "dummy2"), spark.createDataFrame(Seq(("F", "G"), ("H", "I"))).toDF("dummy1", "dummy2") ) // 一行代码完成所有DataFrame的Union操作 // Spark 2.0+版本中union与unionAll功能完全一致,unionAll已被标记为废弃方法 val unionDf = listOfDfs.reduce(_ union _) // 输出结果 unionDf.show()
如果需要处理空List的场景,可以改用foldLeft指定初始空DataFrame:
// 空List场景兼容写法 val emptyDf = spark.emptyDataFrame val unionDf = listOfDfs.foldLeft(emptyDf)(_ union _)
循环写法兼容方案
如果你需要沿用原有的循环逻辑,需要提前在循环外声明可变变量:
import org.apache.spark.sql.DataFrame // 循环外提前声明可变变量var,指定DataFrame类型 var unionDf: DataFrame = null // 遍历所有DataFrame for ((df, i) <- listOfDfs.zipWithIndex) { if (i == 0) { unionDf = df } else { unionDf = unionDf.union(df) } } unionDf.show()
注意事项
- 所有参与合并的DataFrame必须保证列数量、列顺序、列数据类型完全一致,否则会出现数据错位或运行报错
- 如果需要按列名匹配合并,可以先调用
unionByName方法替代union
内容的提问来源于stack exchange,提问作者Error_2646
相关产品推荐
相关产品推荐

