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
相关产品推荐
相关产品推荐

