You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 16:06:29