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

Spark 1.6 SQL中筛选仅含'Continue'状态员工的最优方案咨询

在Spark中筛选仅拥有'Continue'状态员工的最优方法

嘿,我来帮你搞定这个需求——从给定的员工数据里找出**所有记录状态都是'Continue'**的员工,也就是那些从来没有出现过其他状态(比如Finished、No Info)的员工。先看一下你提供的输入数据:

Emp_Id Emp_Age Emp_Status Emp_Name Date_Updated
1 43 Continue John 3/3/15
1 43 Continue John 3/4/15
2 35 Continue Peter 3/5/15
3 32 Finished Alaxender 3/6/15
3 32 Continue Alaxender 3/7/15
4 45 Continue Patrick 3/8/15
4 45 No Info Patrick 3/9/15

核心思路

要实现这个需求,核心是按员工唯一标识(Emp_Id)分组,然后检查该员工的所有状态是否只有'Continue'。下面两种是Spark里效率较高的实现方式:


方法一:利用去重集合判断(直观易读)

这种方式通过collect_set函数收集每个员工的所有状态并去重,然后判断集合是否只有'Continue'这一个元素。代码示例(Scala版本):

import org.apache.spark.sql.functions._

// 假设你已经加载好数据到empDF这个DataFrame中
val resultDF = empDF
  // 按员工唯一信息分组,避免重名干扰
  .groupBy("Emp_Id", "Emp_Name", "Emp_Age")
  .agg(
    collect_set("Emp_Status").alias("status_set"),
    max("Date_Updated").alias("latest_update") // 可选:保留该员工最新的更新日期
  )
  // 筛选出仅含Continue的分组
  .where(size(col("status_set")) === 1 && col("status_set")(0) === "Continue")
  // 移除临时生成的状态集合列
  .drop("status_set")

// 查看结果
resultDF.show()

优点:逻辑非常直观,容易理解和维护,适合中小规模的数据场景。


方法二:利用计数聚合(性能更优)

如果你的数据量很大,用计数类的聚合函数性能会更好——不需要收集整个状态集合,只需要统计不同状态的数量和Continue状态的记录数:

import org.apache.spark.sql.functions._

val resultDF = empDF
  .groupBy("Emp_Id", "Emp_Name", "Emp_Age")
  .agg(
    countDistinct("Emp_Status").alias("distinct_status_count"),
    count(when(col("Emp_Status") === "Continue", true)).alias("continue_record_count"),
    max("Date_Updated").alias("latest_update")
  )
  // 筛选:只有一种状态,且该状态是Continue
  .where(col("distinct_status_count") === 1 && col("continue_record_count") > 0)
  // 移除临时计数列
  .drop("distinct_status_count", "continue_record_count")

resultDF.show()

优点:聚合计算的开销更小,在大数据量场景下性能更出色,适合生产环境的大规模数据处理。


预期输出示例

执行上述代码后,你会得到符合要求的员工数据:

+------+--------+--------+-------------+
|Emp_Id|Emp_Age |Emp_Name|latest_update|
+------+--------+--------+-------------+
|1     |43      |John    |3/4/15       |
|2     |35      |Peter   |3/5/15       |
+------+--------+--------+-------------+

(注:员工4因为有一条"No Info"的记录,所以会被排除在外)

内容的提问来源于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 03:39:11