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

Spark Scala高效实现员工爱好宽表转聚合长表需求问询

Scala Spark高效实现爱好-员工列表映射

问题场景

输入DataFrame:

Emp_IDCricketChessSwim
11YNN
12YYY
13NNY

需要转换为:

HobbiesEmp_id_list
Cricket11,12
Chess12
Swim12,13

之前采用逐个创建爱好DataFrame再合并的方式,现寻求更高效的实现方案。


高效实现方案

核心采用宽表转长表(Unpivot)+ 分组聚合的思路,避免多次拆分合并带来的性能开销,具体步骤如下:

  1. 宽表转长表:用Spark原生stack函数将所有爱好列转换为Hobbies(爱好名称)和Status(是否参与)两列
  2. 过滤无效数据:仅保留Status为Y的记录,减少聚合阶段的数据量
  3. 分组拼接:按爱好分组,收集对应员工ID并拼接为逗号分隔的字符串

Scala 代码示例

import org.apache.spark.sql.functions.{stack, col, concat_ws, collect_list}

// 假设原始DataFrame名为empHobbies
val resultDF = empHobbies
  // 使用stack转换宽表:3代表3个爱好列,每对参数对应"爱好名称", 原列值
  .select(
    col("Emp_ID"),
    stack(3, "Cricket", col("Cricket"), "Chess", col("Chess"), "Swim", col("Swim"))
      .alias("Hobbies", "Status")
  )
  // 过滤出有该爱好的员工记录
  .filter(col("Status") === "Y")
  // 按爱好分组,收集并拼接员工ID(注意转String避免数字类型拼接问题)
  .groupBy("Hobbies")
  .agg(
    concat_ws(",", collect_list(col("Emp_ID").cast("string")))
      .alias("Emp_id_list")
  )

// 查看结果
resultDF.show()

代码说明

  • stack函数:通用的宽表转长表工具,无需额外依赖,参数中的数字对应要转换的列数量,后续按"列名-列值"的成对格式传入即可
  • 提前过滤无效记录:减少聚合阶段处理的数据量,直接提升性能
  • collect_list + concat_ws:高效完成员工ID的收集与拼接,确保输出格式符合要求

这种方案仅需一次数据转换和一次聚合操作,相比逐个拆分合并的方式,大幅减少了Shuffle和IO次数,在爱好列较多的场景下性能优势更明显。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:57:08