PySpark Avro序列化正常,反序列化报Union转Record错误求解
PySpark Avro序列化反序列化问题
在PySpark 3.1.2中使用Avro进行数据序列化与反序列化时,序列化正常,但反序列化报错:
org.apache.spark.sql.avro.IncompatibleSchemaException: Attempting to treat union as a RECORD, but it was: UNION
测试代码
import shutil from pyspark.sql import SparkSession from pyspark.sql.avro.functions import from_avro, to_avro from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 初始化Spark会话 spark = ( SparkSession.builder.appName("AvroSerializationExample") .config("spark.jars.packages", "org.apache.spark:spark-avro_2.12:3.1.2,") .getOrCreate() ) schema = StructType( [ StructField("key", StringType(), True), StructField( "value", StructType( [ StructField("id", IntegerType(), False), StructField("name", StringType(), True), ] ), True, ), ] ) data = [ ("key1", {"id": 1, "name": "Test 1"}), ("key2", {"id": 2, "name": "Test 2"}), ("key3", {"id": 3, "name": "Test 3"}), ("key4", None), (None, None), ] df = spark.createDataFrame(data, schema) # 转换为Avro格式 avro_json_schema = """["null",{ "type": "record", "name": "TestRecord", "namespace": "com.example", "fields": [ {"name": "id", "type": "int"}, {"name": "name", "type": ["null", "string"]} ] }]""" df_with_avro = df.select(to_avro("value", avro_json_schema).alias("value")) # 展示原始和序列化后的数据 print("Original DataFrame:") df.show() print("\nDataFrame with Avro serialized column:") df_with_avro.show(truncate=False) # 写入文件 output_path = "./tmp" shutil.rmtree(output_path, ignore_errors=True) df_with_avro.write.format("parquet").save(output_path) # 读取并反序列化 df_read = spark.read.parquet(output_path) df_deserialized = df_read.select(from_avro("value", avro_json_schema)) df_deserialized.show(truncate=False)
已做测试
- 移除
to_avro的jsonFormatSchema参数后,代码可正常运行,说明to_avro内部自动生成Schema的逻辑与手动指定不同。 - 修改Avro Schema,去掉外层Union包装仅保留Record定义时,反序列化无错误,但会产生警告:
25/04/26 12:03:07 WARN AvroSerializer: Writing Avro files with non-nullable Avro schema and nullable catalyst schema will throw runtime exception if there is a record with null value.,而测试数据包含null值却未报错,与警告内容不符。
疑问
- 该警告是否不准确或具有误导性?
- 为何不含Union的Schema理论上不支持null却能正常运行?
- 在此场景下处理可空Record的正确方式是什么?
问题解答
1. 警告是否不准确或具有误导性?
这个警告不算完全错误,但存在场景适配问题。它主要针对直接写入Avro格式文件的场景:当用df.write.format("avro")写入时,如果Avro Schema是非nullable但Catalyst Schema是nullable,遇到null值会直接抛出异常。但你的场景是用to_avro将列序列化为Avro二进制后写入Parquet,这种情况下to_avro会对null值做特殊处理(比如写入Avro的null类型),不会触发警告里描述的运行时异常,所以警告会显得“不准”,本质是警告触发条件和你的实际场景不匹配。
2. 为何不含Union的Schema理论上不支持null却能正常运行?
Avro规范里,非Union的Schema确实不允许null值,但to_avro函数在处理nullable的Catalyst列时,即使你指定的Avro Schema没有包含null的Union,它也会隐式地将null值序列化为Avro的null类型。而from_avro在读取时,也能兼容这种隐式写入的null值,所以不会报错。这是PySpark Avro库的一个容错处理,但不符合严格的Avro Schema规范,不建议依赖。
3. 处理可空Record的正确方式
正确做法是让Avro Schema和Catalyst Schema严格对齐:
- 因为你的Catalyst Schema中
value字段是nullable(True),所以对应的Avro Schema必须是包含null的Union类型,且Union的第一个元素必须是null(这是之前报错的关键)。 - 手动指定Schema时容易出现结构或元数据不匹配,更稳妥的方式是复用
to_avro自动生成的Schema:
# 序列化时获取自动生成的Avro Schema avro_schema = to_avro("value").schema.simpleString() df_with_avro = df.select(to_avro("value", avro_schema).alias("value")) # 反序列化时使用同一个Schema df_deserialized = df_read.select(from_avro("value", avro_schema))
如果必须手动指定Schema,要确保:
- Union的第一个元素是
null - Record字段和Catalyst Schema完全匹配(包括nullable属性)
- 可以通过Catalyst Schema直接生成对应Avro Schema,保证一致性:
# 从Catalyst Schema生成Avro Schema avro_schema = df.schema["value"].dataType.json() df_with_avro = df.select(to_avro("value", avro_schema).alias("value"))
这样就能避免反序列化的Schema不兼容问题,同时严格遵循Avro规范处理null值。
内容的提问来源于stack exchange,提问作者Manish Zacharias
相关产品推荐
相关产品推荐

