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

Spark 1.6中识别同时含'In Progress'和'Finished'工单的最优方案

Spark 1.6 实现工单状态筛选的最优方案

嘿,针对你这个需求——找出同时存在In Progress和Finished/Not Valid状态的工单,且只保留这些工单的In Progress记录——我给你梳理下最优的实现方案,主要推荐用窗口函数,性能更优,也更简洁:

需求回顾

我们需要:

  1. 识别出那些工单ID(Work_Req_Id)同时包含两种状态:In Progress,以及Finished或Not Valid
  2. 仅保留这些工单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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:03:06