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

如何使用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()

代码说明

  1. 窗口分组:partitionBy("project")确保同一项目的任务被归为一组,orderBy("Task")保证任务行的顺序稳定,确保第一个任务行被标记为行号1。
  2. 行号标记:row_number()函数给每组内的行分配从1开始的递增行号。
  3. 项目证据处理:利用when条件函数,仅保留行号为1的行的Project-Evidence值,其余行该列设为null。
  4. 清理临时列:删掉用于标记的row_num列,得到符合预期的DataFrame。

如果任务的排序规则需要调整,只需修改orderBy中的列即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 20:52:38