如何在Spark Scala DataFrame中合并同ID多行数据为单行
Spark Scala DataFrame 合并同ID多行数据为单行
根据你的需求,以下是几种常见场景下的实现方案:
1. 将同一ID的字段值收集为列表(保留所有/去重值)
适用于需要保留同一ID下所有字段值的场景:
// 导入所需函数 import org.apache.spark.sql.functions.{collect_list, collect_set} // 分组聚合:collect_list保留重复值,collect_set自动去重 val mergedDF = originalDF .groupBy("id") .agg( collect_list("column1").alias("column1_values"), collect_set("column2").alias("column2_unique_values") )
2. 合并互补字段(同一ID下每个字段仅一个非空值)
如果同一ID的多行数据中,各字段的非空值是互补的(比如每行仅填充一个字段),可以用first/last忽略空值合并:
import org.apache.spark.sql.functions.{first, last} val mergedDF = originalDF .groupBy("id") .agg( first("column1", ignoreNulls = true).alias("column1"), last("column2", ignoreNulls = true).alias("column2") )
3. 行转列合并(将同一ID的键值对转为列)
若输入数据是id+key+value的结构,需要将不同key转为单独列:
val pivotedDF = originalDF .groupBy("id") .pivot("key") // 指定要转列的字段 .agg(first("value", ignoreNulls = true))
4. 将字段值拼接为字符串
如果需要把同一ID的字段值拼接成单个字符串,用concat_ws结合collect_list:
import org.apache.spark.sql.functions.{concat_ws, collect_list} val mergedDF = originalDF .groupBy("id") .agg( concat_ws(", ", collect_list("column1")).alias("column1_combined"), concat_ws("|", collect_list("column2")).alias("column2_combined") )
内容的提问来源于stack exchange,提问作者bigdata techie
相关产品推荐
相关产品推荐

