使用PySpark读取Avro文件时自定义Schema转换报错求助
解决PySpark读取Avro自定义Schema时的KeyError: 'nullable'问题
问题场景
尝试通过自定义Schema读取Avro文件,将Avro格式的JSON Schema转换为PySpark的StructType时,触发KeyError: 'nullable'错误。
错误原因
StructType.fromJson()仅支持Spark原生Schema的JSON格式,该格式要求每个字段必须包含nullable属性;而你提供的是Avro原生的Record类型Schema,本身不包含该属性,因此转换时报错。
解决方案
方法1:手动为Avro Schema字段添加nullable属性
Avro中字段默认是可空的(除非字段类型明确排除null),直接给每个字段添加"nullable": true即可:
import json from pyspark.sql.types import StructType from pyspark.sql import SparkSession json_schema = """ { "type": "record", "name": "User", "fields": [ { "name": "routingNumber", "type": "string", "nullable": true } ] } """ schema_dict = json.loads(json_schema) # 提取fields部分,StructType.fromJson期望的是包含fields的结构 avro_schema = StructType.fromJson({"fields": schema_dict["fields"]}) spark = SparkSession.builder.appName("AvroReadExample") .config('spark.jars', '/Users/harbeerkadian/Documents/workspace/learn-spark/jars/spark-avro_2.12-3.5.0.jar') .getOrCreate() df = spark.read.format("avro").schema(avro_schema).load("/Users/harbeerkadian/Downloads/accounts.avro")
方法2:使用Spark-Avro工具直接转换Avro Schema(推荐)
利用Spark-Avro提供的from_avro_schema方法,直接将Avro Schema转为Spark StructType,无需手动修改:
import json from pyspark.sql.avro.functions import from_avro_schema from pyspark.sql import SparkSession json_schema = """ { "type": "record", "name": "User", "fields": [ { "name": "routingNumber", "type": "string" } ] } """ # 直接将Avro Schema字符串转为Spark StructType avro_schema = from_avro_schema(json_schema) spark = SparkSession.builder.appName("AvroReadExample") .config('spark.jars', '/Users/harbeerkadian/Documents/workspace/learn-spark/jars/spark-avro_2.12-3.5.0.jar') .getOrCreate() df = spark.read.format("avro").schema(avro_schema).load("/Users/harbeerkadian/Downloads/accounts.avro")
说明
方法2更符合Avro与Spark集成的规范,尤其是当Schema字段较多或结构复杂时,能避免手动修改的繁琐和出错概率。
内容的提问来源于stack exchange,提问作者Harbeer Kadian
相关产品推荐
相关产品推荐

