Spark Scala 如何将一对列的所有行合并到另一对列中
Spark Scala 实现列对行追加合并
需求说明
将表中一组列的所有行,追加到另一组列的下方,最终仅保留目标合并列。
输入输出示例
输入结构(4列示例)
| 列1 | 列2 | 列3 | 列4 |
|---|---|---|---|
| a1 | b1 | c1 | d1 |
| a2 | b2 | c2 | d2 |
| a3 | b3 | c3 | d3 |
输出结构(合并为2列)
| 列1 | 列2 |
|---|---|
| a1 | b1 |
| a2 | b2 |
| a3 | b3 |
| c1 | d1 |
| c2 | d2 |
| c3 | d3 |
实现方案
核心逻辑为:拆分不同列对,将待合并列对重命名为和目标列对完全一致的列名,再通过union操作合并所有列对数据。
基础实现代码
import org.apache.spark.sql.SparkSession object ColumnMergeDemo { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("ColumnPairMerge") .getOrCreate() // 实际使用时替换为你的数据源读取逻辑,如 spark.read.parquet(路径) import spark.implicits._ val inputDF = Seq( ("a1", "b1", "c1", "d1"), ("a2", "b2", "c2", "d2"), ("a3", "b3", "c3", "d3") ).toDF("col1", "col2", "col3", "col4") // 提取目标列对 val targetPart = inputDF.select("col1", "col2") // 提取待追加列对,重命名为和目标列对一致的列名 val appendPart = inputDF.selectExpr("col3 as col1", "col4 as col2") // 合并两部分数据 val resultDF = targetPart.union(appendPart) // 验证输出/写入结果 resultDF.show() spark.stop() } }
扩展场景处理
- 多列对批量合并:如果有3组及以上列对需要合并,可通过列表批量处理,减少重复代码
// 定义所有需要合并的列对:(原列名1, 原列名2, 输出列名1, 输出列名2) val mergePairs = List( ("col1", "col2", "final_col1", "final_col2"), ("col3", "col4", "final_col1", "final_col2"), ("col5", "col6", "final_col1", "final_col2") ) // 批量生成各部分DF后合并 val resultDF = mergePairs.map{ case (src1, src2, tgt1, tgt2) => inputDF.selectExpr(s"$src1 as $tgt1", s"$src2 as $tgt2") }.reduce(_ union _)
- 保留数据来源:如果需要区分数据来自原表的哪组列,可新增来源标识列
import org.apache.spark.sql.functions.lit val targetPart = inputDF.select("col1", "col2").withColumn("source", lit("原始列1-列2")) val appendPart = inputDF.selectExpr("col3 as col1", "col4 as col2").withColumn("source", lit("追加列3-列4")) val resultDF = targetPart.union(appendPart)
注意事项
union操作要求参与合并的两个DataFrame的列数量、列类型、列顺序完全一致,重命名列时请注意对齐- 如果原数据存在空值,合并时会保留,无需额外处理
内容的提问来源于stack exchange,提问作者Simon Rex
相关产品推荐
相关产品推荐

