Spark 1.6中识别同时含'In Progress'和'Finished'工单的最优方案
Spark 1.6 实现工单状态筛选的最优方案
嘿,针对你这个需求——找出同时存在In Progress和Finished/Not Valid状态的工单,且只保留这些工单的In Progress记录——我给你梳理下最优的实现方案,主要推荐用窗口函数,性能更优,也更简洁:
需求回顾
我们需要:
- 识别出那些工单ID(
Work_Req_Id)同时包含两种状态:In Progress,以及Finished或Not Valid - 仅保留这些工单ID对应的
In Progress状态记录
假设我们的原始DataFrame名为workOrdersDF,结构如下:
+-----------+-----------+-------+-----------+ |Work_Req_Id|Assigned to| Date| Status| +-----------+-----------+-------+-----------+ | R1| John| 3/4/15|In Progress| | R1| George| 3/5/15|In Progress| | R2| Peter|3/6/15|In Progress| | R1| John| 3/7/15| Finished| | R3| Mary|3/8/15| Not Valid| +-----------+-----------+-------+-----------+
方案一:窗口函数(最优)
利用Spark 1.6支持的窗口函数,按工单ID分组后收集所有状态,再基于状态集合过滤:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义窗口:按Work_Req_Id分组,对每个工单的所有记录聚合 val workReqWindow = Window.partitionBy("Work_Req_Id") // 第一步:给每条记录添加工单的所有状态集合 val withStatusSetDF = workOrdersDF.withColumn( "all_statuses", collect_set("Status").over(workReqWindow) // 收集该工单的所有唯一状态 ) // 第二步:过滤出符合条件的记录,再保留In Progress状态 val resultDF = withStatusSetDF .filter( // 同时包含In Progress,且包含Finished或Not Valid array_contains(col("all_statuses"), "In Progress") && (array_contains(col("all_statuses"), "Finished") || array_contains(col("all_statuses"), "Not Valid")) ) .filter(col("Status") === "In Progress") // 只保留In Progress的记录 .drop("all_statuses") // 移除临时辅助列 // 查看结果 resultDF.show()
为什么这个方案最优?
窗口函数只需要一次shuffle(按Work_Req_Id分区),就能完成状态集合的收集,后续过滤都是本地操作,在数据量较大时性能优势非常明显。
方案二:自关联(备选)
如果对窗口函数不熟悉,也可以用自关联的方式实现,但性能稍差(需要多次shuffle):
// 第一步:提取所有包含Finished/Not Valid状态的工单ID val finishedOrInvalidReqs = workOrdersDF .filter(col("Status") === "Finished" || col("Status") === "Not Valid") .select("Work_Req_Id").distinct() // 第二步:提取所有包含In Progress状态的工单ID,并和上面的ID取交集 val targetWorkReqs = workOrdersDF .filter(col("Status") === "In Progress") .select("Work_Req_Id").distinct() .join(finishedOrInvalidReqs, Seq("Work_Req_Id"), "inner") // 第三步:关联原数据,筛选出目标记录 val resultDF = workOrdersDF .join(targetWorkReqs, Seq("Work_Req_Id"), "inner") .filter(col("Status") === "In Progress") // 查看结果 resultDF.show()
注意事项
- Spark 1.6中
collect_set作为窗口函数是支持的,但要注意它会去重状态,如果同一个工单多次出现同一状态,只会保留一个,这对我们的需求没有影响 - 如果你的日期字段需要参与逻辑(比如取最早/最晚的In Progress记录),可以在窗口函数里添加
orderBy和row_number()进一步筛选
内容的提问来源于stack exchange,提问作者Ansip
相关产品推荐
相关产品推荐

