如何在Spark DataFrame中过滤掉county列值仅出现一次的行?
如何过滤Spark DataFrame中county列值仅出现一次的行
错误原因
你直接在filter中使用聚合函数count(county)是不合法的:filter是行级过滤操作,只能基于当前行的字段值判断,而聚合函数(如count)是针对分组/全局的统计计算,无法直接在where/filter子句中使用,这就是报错提示「Aggregate expressions are not valid in where clause」的原因。
解决方案
下面提供两种常用的实现方式:
方法1:使用窗口函数(推荐,无需额外Join)
通过窗口函数计算每个county的出现次数,再过滤掉次数为1的行:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.count // 定义按county分组的窗口 val countyWindow = Window.partitionBy("county") // 新增county出现次数列,过滤后移除该列(可选) val filteredData = fault_data .withColumn("county_occurrences", count("*").over(countyWindow)) .filter("county_occurrences > 1") .drop("county_occurrences") filteredData.show()
方法2:分组统计后关联原表
先筛选出出现次数>1的county列表,再通过Inner Join保留原表中符合条件的行:
// 第一步:统计并筛选出有效county(出现次数>1) val validCounties = fault_data .groupBy("county") .count() .filter("count > 1") .select("county") // 第二步:关联原表,仅保留有效county的行 val filteredData = fault_data.join(validCounties, Seq("county"), "inner") filteredData.show()
如果validCounties数据量很小,可以使用广播Join优化性能:
import org.apache.spark.sql.functions.broadcast val filteredData = fault_data.join(broadcast(validCounties), Seq("county"), "inner")
内容的提问来源于stack exchange,提问作者uint8_t
相关产品推荐
相关产品推荐

