Spark读取Parquet嵌套Schema时重复列报错求助
Spark默认对列名大小写不敏感,会将所有列名统一转为小写进行解析匹配。你的Parquet数据中MODEVALID和ModeValid_1在小写后虽为modevalid和modevalid_1,但开启mergeSchema=true合并不同文件Schema时,若存在仅大小写不同的同名列(比如部分文件是MODEVALID,部分是ModeValid),就会被识别为重复列抛出错误。另外,将DynamoDB的动态M类型(键值对Map)解析为固定结构的StructType也是核心问题——动态列会被当成Struct的固定字段,一旦字段名大小写差异被忽略,就会触发重复冲突。
1. 正确开启大小写敏感模式
注意:Spark的大小写敏感是全局配置,并非读Parquet时的临时option,需要在创建SparkSession时指定:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("DiscountDataProcessing") \ .config("spark.sql.caseSensitive", "true") \ .getOrCreate()
开启后,Spark会严格区分列名的大小写,MODEVALID和ModeValid_1会被视为完全不同的列,不会触发重复错误。
2. 将动态M类型解析为MapType而非StructType
DynamoDB的M是动态键值对结构,适合用Spark的MapType解析,而非固定结构的StructType,可彻底避免动态列引发的Schema冲突。修改你的Schema定义:
from pyspark.sql.types import StructType, StructField, StringType, MapType # 定义DynamoDB属性值的Schema(对应N/S/P等类型) dynamo_value_schema = StructType([ StructField("S", StringType(), nullable=True), StructField("N", StringType(), nullable=True), StructField("P", StringType(), nullable=True) ]) # 重新定义整体Schema schema_disscountcoderule = StructType([ StructField("dynamodb", StructType([ StructField("ApproximateCreationDateTime", StringType(), nullable=True), StructField("NewImage", StructType([ StructField("AccountEmail", dynamo_value_schema, nullable=True), StructField("DiscountData", StructType([ # 将M字段改为MapType,兼容动态键 StructField("M", MapType(StringType(), dynamo_value_schema), nullable=True) ]), nullable=True) ]), nullable=True) ]), nullable=True) ])
这样不管M里有多少动态键(包括MODEVALID、ModeValid_1这类大小写差异的键),都会被作为Map的key存储,不会再触发重复列错误。
3. 调整mergeSchema的使用场景
如果你的Parquet文件仅M内的动态键不同,顶层结构一致,可考虑关闭mergeSchema,配合MapType解析即可,减少不必要的Schema合并复杂度。
修改配置和Schema后,重新读取数据并验证:
df_disscountcoderule = spark.read \ .schema(schema_disscountcoderule) \ .option("recursiveFileLookup", "true") \ .parquet(f"s3://{source_bucket}/{source_discountcodesrules_prefix}") \ .filter(input_file_name().endswith(".parquet")) \ .filter(~input_file_name().contains("processing-failed")) # 查看Schema确认动态列被正确解析为Map df_disscountcoderule.printSchema()
内容的提问来源于stack exchange,提问作者Johanna

