【问题标题】:Spark sql converts bad records to Null while doing datatype castingSpark sql 在进行数据类型转换时将坏记录转换为 Null
【发布时间】:2022-01-03 18:45:42
【问题描述】:

我有以下数据框:

val simpleData = Seq(Row("James ","","Smith","36636","M",3000),
  Row("Michael ","Rose","","40288","M",4000),
  Row("Robert ","","Williams","42114","M",4000),
  Row("Maria ","Anne","Jones","39192","F",4000),
  Row("Jen","Mary","Brown","bad","F",-1)
)
    
val simpleSchema = StructType(Array(
  StructField("firstname",StringType,true),
  StructField("middlename",StringType,true),
  StructField("lastname",StringType,true),
  StructField("id", StringType, true),
  StructField("gender", StringType, true),
  StructField("salary", IntegerType, true)
))
    
val df = spark.createDataFrame(spark.sparkContext.parallelize(simpleData),simpleSchema)

+---------+----------+--------+-----+------+------+
|firstname|middlename|lastname|   id|gender|salary|
+---------+----------+--------+-----+------+------+
|   James |          |   Smith|36636|     M|  3000|
| Michael |      Rose|        |40288|     M|  4000|
|  Robert |          |Williams|42114|     M|  4000|
|   Maria |      Anne|   Jones|39192|     F|  4000|
|      Jen|      Mary|   Brown|Rose |     F|    -1|
+---------+----------+--------+-----+------+------+

我在下面运行示例代码,我想在转换后将字符串列转换为整数。

df.createOrReplaceTempView("EMP")
val df2 = spark.sql("select cast(id as INT) from EMP")

+-----+
|   id|
+-----+
|36636|
|40288|
|42114|
|39192|
| null|
+-----+

这里所有整数数据都正确转换,但“Rose”转换为空。

能否请您帮我了解如何在出现不良记录时抛出异常? 是否有任何火花配置设置?

此外,如果查询中有多个强制转换,如何获取存在此问题的确切列名。

【问题讨论】:

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


    【解决方案1】:

    由于 Spark 3.0 和票证 SPARK-30292 的更正,当您尝试将无效字符串转换为数字时,将 spark.sql.ansi.enabled 配置设置为 true 将引发异常:

    spark.conf.set("spark.sql.ansi.enabled", "true")
    df.createOrReplaceTempView("EMP")
    val df2 = spark.sql("select cast(id as INT) from EMP")
    

    抛出NumberFormatException。详情请见https://spark.apache.org/docs/latest/sql-ref-ansi-compliance.html#cast

    【讨论】:

    • 谢谢,@Vincent 这个解决方案已经奏效。如果可能的话,你能帮我解决第二个问题吗?如果有多个强制转换,那么很难在错误日志中找到确切的列名,我看不到任何生成异常的列名。
    • 您好,很遗憾,我不知道如何回答您问题的第二部分。也许正如Robert Kossendey's answer 中所建议的那样,通过创建一个UDF
    【解决方案2】:

    如果转换出错,Spark 不会抛出。

    作为捕获这些错误的自定义方法,您可以编写一个UDF,如果您强制转换为 null,则会抛出该错误。这会降低脚本的性能,因为 Spark 无法优化 UDF 执行。

    【讨论】:

      猜你喜欢
      • 2012-03-10
      • 2010-09-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-07-13
      • 2021-04-26
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多