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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 19:57:22