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

Spark中按复杂条件合并DataFrame行(优先无UDF实现)

问题描述

原始Spark DataFrame定义

val df= Seq(
  (100020,"GUID1","CreateWorkflow", "Workflow_Read", "None",  "10-02-2018" ),
  (100020,"GUID1","None","WorkItem_StateChange", "INPROGRESS", "12-02-2018"),
  (100020,"GUID1","None","WorkItem_StateChange", "INPROGRESS", "11-02-2018"),
  (100057,"GUID2","CreateWorkflow", "Workflow_Read", "None",  "10-05-2021" ),
  (100057,"GUID2","None","WorkItem_StateChange", "INPROGRESS", "12-05-2021"),
  (100057,"GUID2","None","WorkItem_StateChange", "INPROGRESS", "11-05-2021")
).toDF("SerialNumber","GUID","UseCase","EventName","NewState", "Time")

原始数据展示

+------------+-----+--------------+--------------------+----------+----------+
|SerialNumber|GUID |UseCase       |EventName           |NewState  |Time      |
+------------+-----+--------------+--------------------+----------+----------+
|100020      |GUID1|CreateWorkflow|Workflow_Read       |None      |10-02-2018|
|100020      |GUID1|None          |WorkItem_StateChange|INPROGRESS|12-02-2018|
|100020      |GUID1|None          |WorkItem_StateChange|INPROGRESS|11-02-2018|
|100057      |GUID2|CreateWorkflow|Workflow_Read       |None      |10-05-2021|
|100057      |GUID2|None          |WorkItem_StateChange|INPROGRESS|12-05-2021|
|100057      |GUID2|None          |WorkItem_StateChange|INPROGRESS|11-05-2021|
+------------+-----+--------------+--------------------+----------+----------+

需求说明

按SerialNumber和GUID分组,将每组内的Workflow_Read行与时间最近的WorkItem_StateChange行合并,剩余的WorkItem_StateChange行保留为独立行。合并后的行需将原两行的时间分别存入CreateTime和InProgresstime列,最终期望结果如下:

期望结果

+------------+-----+--------------+----------+----------+--------------+
|SerialNumber|GUID |UseCase       |NewState  |CreateTime|InProgresstime|
+------------+-----+--------------+----------+----------+--------------+
|100020      |GUID1|CreateWorkflow|INPROGRESS|10-02-2018|11-02-2018    |
|100020      |GUID1|None          |INPROGRESS|None      |12-02-2018    |
|100057      |GUID2|CreateWorkflow|INPROGRESS|10-05-2021|11-05-2021    |
|100057      |GUID2|None          |INPROGRESS|None      |12-05-2021    |
+------------+-----+--------------+----------+----------+--------------+

无UDF高效实现方案

核心思路

通过拆分数据集、窗口函数标记目标行、关联合并三个步骤实现,全程使用Spark内置算子,避免UDF带来的性能损耗。

代码实现

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.DateType
import org.apache.spark.sql.Window

// 1. 转换时间列为日期类型,方便后续排序比较
val df_with_date = df.withColumn("date_time", to_date(col("Time"), "dd-MM-yyyy"))

// 2. 拆分两类数据:创建流程行、状态变更行
val df_create = df_with_date.filter(col("EventName") === "Workflow_Read").alias("c")
val df_progress = df_with_date.filter(col("EventName") === "WorkItem_StateChange").alias("p")

// 3. 用窗口函数给每组状态变更行按时间升序排名,标记出时间最早的行(与创建行最近)
val windowSpec = Window.partitionBy("SerialNumber", "GUID").orderBy("date_time")
val df_progress_ranked = df_progress.withColumn("rank", row_number().over(windowSpec))

// 4. 关联创建行与排名第一的状态变更行,生成合并后的记录
val merged_create = df_create.join(
  df_progress_ranked.filter(col("rank") === 1).alias("p_top"),
  Seq("SerialNumber", "GUID"),
  "inner"
).select(
  col("SerialNumber"),
  col("GUID"),
  $"c.UseCase",
  $"p_top.NewState",
  $"c.Time".alias("CreateTime"),
  $"p_top.Time".alias("InProgresstime")
)

// 5. 处理剩余的状态变更行(排名>1的行),补充CreateTime为None
val remaining_progress = df_progress_ranked.filter(col("rank") > 1)
  .select(
    col("SerialNumber"),
    col("GUID"),
    col("UseCase"),
    col("NewState"),
    lit("None").alias("CreateTime"),
    col("Time").alias("InProgresstime")
  )

// 6. 合并所有结果并排序
val final_df = merged_create.union(remaining_progress)
  .orderBy("SerialNumber", "GUID", "CreateTime")

// 查看最终结果
final_df.show()

性能说明

  • 拆分数据集后分别处理,减少关联操作的数据量,提升效率;
  • 使用row_number()窗口函数精准定位目标行,逻辑清晰且执行高效;
  • 全程依赖Spark内置函数,避免UDF带来的序列化/反序列化开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 18:35:18