【问题标题】:Check count of a column from a dataframe and and add column and count as Map检查数据框中列的计数,并添加列并计数为 Map
【发布时间】:2020-01-14 20:24:18
【问题描述】:

我是 Scala 初学者。我正在尝试在表的列中查找空值的计数,并将列名和计数添加为 Map 中的键值对。下面的代码不能按预期工作。请指导我如何修改此代码以使其正常工作

def nullCheck(databaseName:String,tableName:String) ={
var map = scala.collection.mutable.Map[String, Int]() 
validationColumn = Array(col1,col2)
for(i <- 0 to validationColumn.length) {
val nullVal = spark.sql(s"select count(*) from $databaseName.$tableName where validationColumn(i) is NULL")
if(nullval == 0)
map(validationColumn(i)) = nullVal
map
} 

函数应该返回 ((col1,count),(col2,count)) 作为 Map

【问题讨论】:

    标签: scala apache-spark hadoop hive apache-spark-sql


    【解决方案1】:

    这可以通过创建一个动态 sql 字符串然后映射它来完成。您的方法多次读取相同的数据

    这里是解决方案。我使用了一个“示例” DataFrame。

    scala> val inputDf = Seq((Some("Sam"),None,200),(None,Some(31),30),(Some("John"),Some(25),25),(Some("Harry"),None,100)).toDF("name","age","not_imp_column")
    
    scala> inputDf.show(false)
    +-----+----+--------------+
    |name |age |not_imp_column|
    +-----+----+--------------+
    |Sam  |null|200           |
    |null |31  |30            |
    |John |25  |25            |
    |Harry|null|100           |
    +-----+----+--------------+
    

    我们的 ValidationColumns 是 nameage,我们将在其中计算 Null 我们把它们放在一个列表中

    scala&gt; val validationColumns = List("name","age")

    我们创建一个 SQL 字符串,它将驱动整个计算

    scala&gt; val sqlStr = "select " + validationColumns.map(x =&gt; "sum(" + x + "_count) AS " + x + "_sum" ).mkString(",") + " from (select " + validationColumns.map(x =&gt; "case when " + x + " = '$$' then 1 else 0 end AS " + x + "_count").mkString(",") + " from " +" (select" + validationColumns.map(x =&gt; " nvl( " + x +",'$$') as " + x).mkString(",") + " from example_table where " + validationColumns.map(x =&gt; x + " is null ").mkString("or ") + " ) layer1 ) layer2 "

    解析为 ==>

    "select sum(name_count) AS name_sum,sum(age_count) AS age_sum from (select case when name = '$$' then 1 else 0 end AS name_count,case when age = '$$' then 1 else 0 end AS age_count from (select nvl( name,'$$') as name, nvl( age,'$$') as age from example_table where name is null or age is null ) layer1 ) layer2 "

    现在我们为我们的数据框创建一个临时视图

    inputDf.createOrReplaceTempView("example_table")

    唯一要做的事情就是执行sql并创建一个Map,由

    validationColumns zip spark.sql(sqlStr).collect.map(_.toSeq).flatten.toList toMap

    结果

    Map(name -&gt; 1, age -&gt; 2) // 显然你可以让它类型安全

    【讨论】:

    • 谢谢!我尝试了上述方法。但我收到一个错误。 Spark 是否支持别名? org.apache.spark.sql.catalyst.parser.ParseException:不匹配的输入'来自'期望 (在 org.apache.spark.sql.catalyst.parser.ParseException.withCommand(ParseDriver.scala:239) 在 org. org.apache.spark.sql.execution.SparkSqlParser.parse(SparkSqlParser.scala:48) 在 org.apache.spark.sql 的 apache.spark.sql.catalyst.parser.AbstractSqlParser.parse(ParseDriver.scala:115)。 catalyst.parser.AbstractSqlParser.parsePlan(ParseDriver.scala:69)
    • 你能粘贴 1) 版本的 Spark + 2) 你特别尝试过的 Sql 构建部分的代码吗?
    • 1. 2.3.0.cloudera4 2.val sqlStr = "select " + nullValidationColumn.map(x =&gt; "sum(" + x + "_null_count) AS " + x + "_null_count" ).mkString(",") + " from (select " + nullValidationColumn.map(x =&gt; "case when " + x + " = '$$' then 1 else 0 end AS " + x + "_null_count").mkString(",") + " from" +" (select" + nullValidationColumn.map(x =&gt; " nvl( " + x +",'$$') as " + x).mkString(",") + " from driver_vw where " + nullValidationColumn.map(x =&gt; x + " is null ").mkString("or ") + ") inner ) outer"
    • @vvazza 你能做一个 spark.sql(“select * from driver_vw”).show(false) 它显示什么?
    • ` |c1_sk|c2_sk| |201372784 |2013727841 | |201372800 |2013728002 | |201372792 |2013727923 | |201372791 |2013727914 | |199795779 |1997957795 | `
    猜你喜欢
    • 2016-02-02
    • 2019-04-19
    • 1970-01-01
    • 1970-01-01
    • 2020-03-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多