Spark SQL类型转换时坏记录转NULL,如何抛出异常并定位问题列?
问题1:类型转换遇到坏记录时抛出异常的实现方式
Spark 3.0及以上版本提供了标准配置项实现该需求:
- 开启ANSI SQL模式配置:
spark.sql.ansi.enabled = true
开启该配置后,所有cast操作遇到不符合目标类型的输入时,会直接抛出IllegalArgumentException异常终止任务执行,而非返回Null。
示例代码:
// 开启ANSI模式 spark.conf.set("spark.sql.ansi.enabled", "true") df.createOrReplaceTempView("EMP") // 此时执行该SQL会直接抛出类型转换失败异常 val df2 = spark.sql("select cast(id as INT) from EMP")
如果是Spark 2.x版本,没有原生ANSI配置支持,可以自定义UDF实现类型转换,在UDF逻辑中判断转换失败直接抛出异常即可。
问题2:多列类型转换时定位出错列
有两种常用方案可以精准定位出错的列:
方案1:依赖ANSI模式的异常信息
开启spark.sql.ansi.enabled后,抛出的异常信息中会包含失败的转换表达式,例如CAST(id AS INT),直接从异常栈中即可读取到出错的列名为id。
方案2:使用try_cast主动校验(适合不希望终止任务的场景)
try_cast是Spark 2.3及以上版本支持的函数,转换失败时会返回Null,不会抛出异常。可以对每个需要转换的列单独做try_cast校验,对比原值和转换后的值判断是否转换失败,同时标记出问题列:
示例代码:
val checkDf = spark.sql(""" select *, -- 拼接所有转换失败的列名 concat_ws(',', case when id is not null and try_cast(id as INT) is null then 'id' end, case when gender is not null and try_cast(gender as BOOLEAN) is null then 'gender' end ) as error_columns from EMP """) // 过滤出存在转换错误的记录 checkDf.filter("error_columns != ''").show(false)
执行后error_columns列会列出当前行所有转换失败的列名,既可以统计错误情况,也不会中断正常任务执行。
内容的提问来源于stack exchange,提问作者Ranjan Panda
相关产品推荐
相关产品推荐

