You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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值却未报错,与警告内容不符。

疑问

  1. 该警告是否不准确或具有误导性?
  2. 为何不含Union的Schema理论上不支持null却能正常运行?
  3. 在此场景下处理可空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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 08:29:55