Spark SQL分组与字符串聚合:如何转换DataFrame结构?
Spark SQL实现DataFrame分组与字符串聚合
根据你的需求,核心是按指定列分组后,对目标列做字符串拼接聚合,下面提供两种常用实现方式:
一、DataFrame API 方式
Scala 版本
首先导入Spark内置的聚合函数:
import org.apache.spark.sql.functions.{collect_list, concat_ws}
然后编写分组聚合逻辑:
// originalDF 是你的原始DataFrame val resultDF = originalDF .groupBy("id", "name") // 按id和name分组 .agg( // 把subject列的所有值用逗号拼接成字符串,命名为subjects concat_ws(",", collect_list("subject")).alias("subjects"), // 把score列的所有值用逗号拼接成字符串,命名为scores concat_ws(",", collect_list("score")).alias("scores") )
Python 版本
导入函数后执行相同逻辑:
from pyspark.sql.functions import collect_list, concat_ws result_df = original_df \ .groupBy("id", "name") \ .agg( concat_ws(",", collect_list("subject")).alias("subjects"), concat_ws(",", collect_list("score")).alias("scores") )
二、Spark SQL 语句方式
先把原始DataFrame注册为临时视图:
// Scala/Python 都适用 originalDF.createOrReplaceTempView("student_scores")
然后执行SQL查询:
SELECT id, name, concat_ws(',', collect_list(subject)) AS subjects, -- 如果score是数值类型,需要先转成字符串再拼接 concat_ws(',', collect_list(cast(score as string))) AS scores FROM student_scores GROUP BY id, name
补充说明
collect_list:会保留分组内所有重复的列值,如果需要去重聚合,可以换成collect_setconcat_ws:第一个参数是分隔符,第二个参数是要拼接的列表,能自动处理空值情况
内容的提问来源于stack exchange,提问作者Trowa
相关产品推荐
相关产品推荐

