如何使用Scala在Spark DataFrame上应用自定义逻辑(嵌套JSON场景)
用Scala处理扁平化后DataFrame的自定义逻辑实现
问题背景
我有嵌套JSON数据,已通过explode()扁平化得到包含project、Task、Task-Evidence、Task-Remarks、Project-Evidence列的DataFrame。其中:
- 包含1个项目,对应2个任务
- 两个任务各有1个任务链接
- 项目层级有3个项目链接
当前DataFrame的问题是:每个任务对应的行都重复显示了全部3个项目证据;预期要实现的效果是:仅第一个任务所在的行显示项目证据,第二个任务所在行的项目证据列留空;两个任务的Task、Task-Evidence、Task-Remarks列分别对应各自的任务数据,project列保持一致。
实现方案
可以通过Spark的窗口函数(Window)标记分组内的行号,再根据行号决定是否保留Project-Evidence值,具体Scala代码如下:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 定义窗口:按project分组,给同项目下的任务行分配行号 val projectWindow = Window.partitionBy("project").orderBy("Task") // 处理DataFrame:标记行号,按需保留项目证据 val resultDF = originalDF .withColumn("row_num", row_number().over(projectWindow)) .withColumn("Project-Evidence", when(col("row_num") === 1, col("Project-Evidence")).otherwise(null)) .drop("row_num") // 移除临时行号列 // 查看最终结果 resultDF.show()
代码说明
- 窗口分组:
partitionBy("project")确保同一项目的任务被归为一组,orderBy("Task")保证任务行的顺序稳定,确保第一个任务行被标记为行号1。 - 行号标记:
row_number()函数给每组内的行分配从1开始的递增行号。 - 项目证据处理:利用
when条件函数,仅保留行号为1的行的Project-Evidence值,其余行该列设为null。 - 清理临时列:删掉用于标记的
row_num列,得到符合预期的DataFrame。
如果任务的排序规则需要调整,只需修改orderBy中的列即可。
内容的提问来源于stack exchange,提问作者Vivek Gowda
相关产品推荐
相关产品推荐

