Spark DataFrame读取JSON时字段类型不匹配导致整行返回NULL问题排查
问题描述
输入数据
{"driverId":1,"driverRef":"hamilton","number":44,"code":"HAM","name":{"forename":"Lewis","surname":"Hamilton"},"dob":"1985-01-07","nationality":"British","url":"http://en.wikipedia.org/wiki/Lewis_Hamilton"} {"driverId":2,"driverRef":"heidfeld","number":"\N","code":"HEI","name":{"forename":"Nick","surname":"Heidfeld"},"dob":"1977-05-10","nationality":"German","url":"http://en.wikipedia.org/wiki/Nick_Heidfeld"} {"driverId":3,"driverRef":"rosberg","number":6,"code":"ROS","name":{"forename":"Nico","surname":"Rosberg"},"dob":"1985-06-27","nationality":"German","url":"http://en.wikipedia.org/wiki/Nico_Rosberg"} {"driverId":4,"driverRef":"alonso","number":14,"code":"ALO","name":{"forename":"Fernando","surname":"Alonso"},"dob":"1981-07-29","nationality":"Spanish","url":"http://en.wikipedia.org/wiki/Fernando_Alonso"} {"driverId":5,"driverRef":"kovalainen","number":"\N","code":"KOV","name":{"forename":"Heikki","surname":"Kovalainen"},"dob":"1981-10-19","nationality":"Finnish","url":"http://en.wikipedia.org/wiki/Heikki_Kovalainen"} {"driverId":6,"driverRef":"nakajima","number":"\N","code":"NAK","name":{"forename":"Kazuki","surname":"Nakajima"},"dob":"1985-01-11","nationality":"Japanese","url":"http://en.wikipedia.org/wiki/Kazuki_Nakajima"} {"driverId":7,"driverRef":"bourdais","number":"\N","code":"BOU","name":{"forename":"Sébastien","surname":"Bourdais"},"dob":"1979-02-28","nationality":"French","url":"http://en.wikipedia.org/wiki/S%C3%A9bastien_Bourdais"}
将上述数据读取为Spark DataFrame后展示时,driverId为2、5、6、7的行全部显示为NULL,对应行的number字段值为\N。使用的代码如下:
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, DateType name_field = StructType(fields =[ StructField("forename", StringType(), True), StructField("surname", StringType(), True) ]) driver_schema = StructType(fields =[ StructField("driverId", IntegerType(), False), StructField("driverRef", StringType(), True), StructField("number", IntegerType(), True), StructField("code", StringType(), True), StructField("name", name_field), StructField("dob", DateType(), True), StructField("nationality", StringType(),True), StructField("url", StringType(), True) ]) driver_df = spark.read\ .schema(driver_schema)\ .json('dbfs:/mnt/databrickslearnf1azure/raw/drivers.json') driver_df.printSchema() root |-- driverId: integer (nullable = true) |-- driverRef: string (nullable = true) |-- number: integer (nullable = true) |-- code: string (nullable = true) |-- name: struct (nullable = true) | |-- forename: string (nullable = true) | |-- surname: string (nullable = true) |-- dob: date (nullable = true) |-- nationality: string (nullable = true) |-- url: string (nullable = true) display(driver_df)
异常结果展示:
问题解答
错误原因
你在schema中将number字段指定为IntegerType,但输入数据里部分行的number取值为字符串\N,和声明的类型不匹配。Spark读取JSON默认使用PERMISSIVE模式,遇到类型不匹配的行时会将整行所有字段设为NULL,这就是部分行全为NULL的根本原因。
修复方案
方案1:读取时指定空值标识
直接在读取JSON时添加nullValue参数,让Spark自动将\N识别为空值,无需修改原有schema定义:
driver_df = spark.read\ .schema(driver_schema)\ .option("nullValue", "\\N")\ .json('dbfs:/mnt/databrickslearnf1azure/raw/drivers.json')
方案2:调整字段类型后转换
先将number字段的schema声明改为StringType,读取完成后再做类型转换,将\N替换为空值:
from pyspark.sql.functions import col, when, trim # 修改number字段类型为字符串 driver_schema = StructType(fields =[ StructField("driverId", IntegerType(), False), StructField("driverRef", StringType(), True), StructField("number", StringType(), True), StructField("code", StringType(), True), StructField("name", name_field), StructField("dob", DateType(), True), StructField("nationality", StringType(),True), StructField("url", StringType(), True) ]) driver_df = spark.read.schema(driver_schema).json('dbfs:/mnt/databrickslearnf1azure/raw/drivers.json') # 转换为整数类型,\N替换为空 driver_df = driver_df.withColumn("number", when(trim(col("number")) == "\\N", None).otherwise(col("number").cast(IntegerType())))
内容的提问来源于stack exchange,提问作者pcbzmani
相关产品推荐
相关产品推荐

