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

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_set
  • concat_ws:第一个参数是分隔符,第二个参数是要拼接的列表,能自动处理空值情况

内容的提问来源于stack exchange,提问作者Trowa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 18:41:00