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

Spark Scala计算列唯一值占总行数百分比时selectExpr类型不匹配问题排查

问题成因
  1. 接口参数不匹配:Spark DataFrame的selectExpr方法要求入参为多个SQL表达式字符串(类型为String*),你传入的是Iterable[(String, Long, Long)]格式的元组集合,类型完全不符合接口要求,这是编译报错的直接原因。
  2. 业务逻辑错误:原有代码中data.head().getValuesMap[Long](data.columns)仅能读取DataFrame第一行的各列值,无法统计每列的全局唯一值数量,即使解决类型错误也无法得到预期结果。
修复方案

首先明确需求是统计每列的唯一值数、及该值占总行数的百分比,可按以下两种场景实现:

场景1:返回(列名, 唯一值数, 占比)格式的本地集合

适合小数据量下快速遍历判断要删除的列:

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

def uniqueValuesAsPercentage(data: org.apache.spark.sql.DataFrame) = {
  val rowCount = data.count()
  // 生成所有列的聚合表达式
  val aggExpressions = data.columns.flatMap(colName => Seq(
    countDistinct(col(colName)).as(s"${colName}_cnt"),
    round(countDistinct(col(colName)) / lit(rowCount) * 100, 2).as(s"${colName}_pct")
  ))
  // 执行一次聚合计算所有列的统计值
  val aggResult = data.agg(aggExpressions.head, aggExpressions.tail: _*).head()
  // 转换为指定格式的集合
  data.columns.map(colName => {
    val uniqueCnt = aggResult.getAs[Long](s"${colName}_cnt")
    val uniquePct = aggResult.getAs[Double](s"${colName}_pct")
    (colName, uniqueCnt, uniquePct)
  })
}

场景2:返回DataFrame格式的统计结果

适合大数据量下直接用SQL/DataSet API过滤要删除的列:

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

def uniqueValuesAsPercentageDf(data: org.apache.spark.sql.DataFrame) = {
  val rowCount = data.count()
  val statRows = data.columns.map(colName => {
    val uniqueCnt = data.select(countDistinct(col(colName))).head().getLong(0)
    val uniquePct = uniqueCnt * 100.0 / rowCount
    Row(colName, uniqueCnt, uniquePct)
  })
  val statSchema = StructType(Seq(
    StructField("column_name", StringType),
    StructField("unique_count", LongType),
    StructField("unique_percentage", DoubleType)
  ))
  data.sparkSession.createDataFrame(data.sparkSession.sparkContext.parallelize(statRows), statSchema)
}

使用示例

过滤删除唯一值占比超过95%的列:

val statDf = uniqueValuesAsPercentageDf(originalDf)
val columnsToDrop = statDf.filter(col("unique_percentage") > 95).select("column_name").as[String].collect()
val cleanedDf = originalDf.drop(columnsToDrop: _*)

内容的提问来源于stack exchange,提问作者joesan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 04:36:07