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

