Spark Scala高效实现员工爱好宽表转聚合长表需求问询
Scala Spark高效实现爱好-员工列表映射
问题场景
输入DataFrame:
| Emp_ID | Cricket | Chess | Swim |
|---|---|---|---|
| 11 | Y | N | N |
| 12 | Y | Y | Y |
| 13 | N | N | Y |
需要转换为:
| Hobbies | Emp_id_list |
|---|---|
| Cricket | 11,12 |
| Chess | 12 |
| Swim | 12,13 |
之前采用逐个创建爱好DataFrame再合并的方式,现寻求更高效的实现方案。
高效实现方案
核心采用宽表转长表(Unpivot)+ 分组聚合的思路,避免多次拆分合并带来的性能开销,具体步骤如下:
- 宽表转长表:用Spark原生
stack函数将所有爱好列转换为Hobbies(爱好名称)和Status(是否参与)两列 - 过滤无效数据:仅保留
Status为Y的记录,减少聚合阶段的数据量 - 分组拼接:按爱好分组,收集对应员工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
相关产品推荐
相关产品推荐

