【问题标题】:Spark Scala replace Dataframe blank records to "0"Spark Scala 将 Dataframe 空白记录替换为“0”
【发布时间】:2017-03-14 13:40:10
【问题描述】:

我需要将我的 Dataframe 字段的空白记录替换为“0”

这是我的代码 -->

import sqlContext.implicits._

case class CInspections (business_id:Int, score:String, date:String, type1:String)

val baseDir = "/FileStore/tables/484qrxx21488929011080/"
val raw_inspections = sc.textFile (s"$baseDir/inspections_plus.txt")
val raw_inspectionsmap = raw_inspections.map ( line => line.split ("\t"))
val raw_inspectionsRDD = raw_inspectionsmap.map ( raw_inspections => CInspections (raw_inspections(0).toInt,raw_inspections(1), raw_inspections(2),raw_inspections(3)))
val raw_inspectionsDF = raw_inspectionsRDD.toDF
raw_inspectionsDF.createOrReplaceTempView ("Inspections")
raw_inspectionsDF.printSchema
raw_inspectionsDF.show()

我正在使用案例类,然后转换为 Dataframe。但我需要“分数”作为 Int,因为我必须执行一些操作并对其进行排序。 但是,如果我将其声明为 score:Int,那么我会收到空​​白值错误。

java.lang.NumberFormatException:对于输入字符串:“”

+-----------+-----+--------+--------------------+
|business_id|score|    date|               type1|
+-----------+-----+--------+--------------------+
|         10|     |20140807|Reinspection/Foll...|
|         10|   94|20140729|Routine - Unsched...|
|         10|     |20140124|Reinspection/Foll...|
|         10|   92|20140114|Routine - Unsched...|
|         10|   98|20121114|Routine - Unsched...|
|         10|     |20120920|Reinspection/Foll...|
|         17|     |20140425|Reinspection/Foll...|
+-----------+-----+--------+--------------------+

我需要将 score 字段作为 Int,因为对于以下查询,它排序为 String 而不是 Int 并给出错误结果

sqlContext.sql("""select raw_inspectionsDF.score  from raw_inspectionsDF where score <>"" order by score""").show()

+-----+
|score|
+-----+
|  100|
|  100|
|  100|
+-----+

【问题讨论】:

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


    【解决方案1】:

    空字符串无法转成Integer,需要将Score设为nullable,这样如果字段缺失,则表示为null,可以尝试以下方法:

    import scala.util.{Try, Success, Failure}
    

    1) 定义一个自定义的解析函数,如果字符串不能转换为 Int,则返回 None,在你的情况下为空字符串;

    def parseScore(s: String): Option[Int] = {
      Try(s.toInt) match {
        case Success(x) => Some(x)
        case Failure(x) => None
      }
    }
    

    2) 将案例类中的 score 字段定义为Option[Int] 类型;

    case class CInspections (business_id:Int, score: Option[Int], date:String, type1:String)
    
    val raw_inspections = sc.textFile("test.csv")
    val raw_inspectionsmap = raw_inspections.map(line => line.split("\t"))
    

    3)使用自定义的parseScore函数解析score字段;

    val raw_inspectionsRDD = raw_inspectionsmap.map(raw_inspections => 
        CInspections(raw_inspections(0).toInt, parseScore(raw_inspections(1)), 
                     raw_inspections(2),raw_inspections(3)))
    
    val raw_inspectionsDF = raw_inspectionsRDD.toDF
    raw_inspectionsDF.createOrReplaceTempView ("Inspections")
    
    raw_inspectionsDF.printSchema
    //root
    // |-- business_id: integer (nullable = false)
    // |-- score: integer (nullable = true)
    // |-- date: string (nullable = true)
    // |-- type1: string (nullable = true)
    
    raw_inspectionsDF.show()
    
    +-----------+-----+----+-----+
    |business_id|score|date|type1|
    +-----------+-----+----+-----+
    |          1| null|   a|    b|
    |          2|    3|   s|    k|
    +-----------+-----+----+-----+
    

    4) 正确解析文件后,可以使用na函数fill轻松将空值替换为0:

    raw_inspectionsDF.na.fill(0).show
    +-----------+-----+----+-----+
    |business_id|score|date|type1|
    +-----------+-----+----+-----+
    |          1|    0|   a|    b|
    |          2|    3|   s|    k|
    +-----------+-----+----+-----+
    

    【讨论】:

    • 非常感谢您的及时回复!它现在工作。 :)
    • 我可以在 sqlContext.sql 中编写如下的 sql 查询吗?我收到以下查询错误-> sqlContext.sql("""select CBusinesses.BUSINESS_ID,CBusinesses.name, CBusinesses.address, CBusinesses.city, CBusinesses.postal_code, CBusinesses.latitude, CBusinesses.longitude, Inspections_notnull.score from CBusinesses , Inspections_notnull where Inspections_notnull.score 0 and CBusinesses.BUSINESS_ID=Inspections_notnull.BUSINESS_ID """).show() java.lang.NumberFormatException: For input string: ""
    • 不太清楚答案,但您似乎正在尝试合并两个表,也许您想要一个联接?
    • 是的。我想计算哪 10 家企业得分最低?”我有 2 个表——“企业”和“检查”,带有 business_id 通用键。它适用于 sql,但如果我在 spark 中使用相同的查询,它就不起作用。如何将 2 个表与一列连接并使用 Spark Sql 计算最高分?我也试过 //val df = businessDF.join(raw_inspectionsDF,businesssDF.col("BUSINESS_ID") == raw_inspectionsDF.col("BUSINESS_ID"))但它也给出错误
    猜你喜欢
    • 2010-12-30
    • 2017-01-22
    • 2018-03-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-01-27
    • 1970-01-01
    相关资源
    最近更新 更多