Spark 2.x到3.x迁移后Parquet中Map类型的Schema路径变更导致不匹配的解决方法咨询
大家好,最近我在负责把团队的Spark工作流从2.4版本迁移到3.5版本时,遇到了一个Parquet文件Schema不兼容的问题——Map类型的嵌套字段路径在两个版本中不一样,导致跨版本读写时报Schema不匹配,这里把问题细节和我找到的解决方法分享出来。
问题复现场景
先给大家看一下我用来复现问题的测试代码(Java版本):
SparkSession sparkSession = SparkSession .builder() .master("local[*]") .config("spark.sql.parquet.writeLegacyFormat", true) .getOrCreate(); // 定义包含Map类型的Schema StructType schema = DataTypes.createStructType(new StructField[]{ DataTypes.createStructField("id", DataTypes.IntegerType, true), DataTypes.createStructField("my_map_data", DataTypes.createMapType(DataTypes.StringType, DataTypes.StringType, Boolean.TRUE), true) }); // 构造测试数据 List<Row> listRows = new ArrayList<>(); listRows.add(RowFactory.create(1, new HashMap<String, String>() {{ put("key", "value"); }})); Dataset<Row> dataset = sparkSession.createDataFrame(listRows, schema); // 写入Parquet文件 dataset.write() .mode(SaveMode.Overwrite) .parquet("spark_all_legacy_write_enabled");
我先在Spark 2.4中运行这段代码生成Parquet文件,然后在Spark 3.5中读取或写入相同路径的文件时,发现了Schema路径的核心差异:
- Spark 2.4生成的Map类型路径:
my_map_data.map.key和my_map_data.map.value - Spark 3.5默认生成的Map类型路径:
my_map_data.key_value.key和my_map_data.key_value.value
可以看到,旧版本用的是map作为Map的嵌套父节点,新版本默认用key_value,即使我开启了spark.sql.parquet.writeLegacyFormat=true也没解决这个问题,导致跨版本读写时Schema校验失败。
问题根源
查了Spark的迁移文档才知道,Spark 3.0+对Parquet格式中Map类型的默认嵌套字段名做了变更:从旧的map改成了key_value,而spark.sql.parquet.writeLegacyFormat参数主要是兼容其他旧的Parquet写入行为(比如Timestamp类型、Decimal类型的处理),但并不覆盖Map类型的字段名变更,所以需要额外配置专门的参数来兼容。
解决方法
我试了几个方案,最终有效的是在Spark 3.x的配置中添加**spark.sql.parquet.legacyMapFieldName=true**参数,具体实现如下:
1. 配置SparkSession实现兼容
在Spark 3.5中读写Parquet时,同时开启两个Legacy参数:
SparkSession sparkSession = SparkSession .builder() .master("local[*]") // 原有Legacy写入格式参数,兼容其他旧Parquet行为 .config("spark.sql.parquet.writeLegacyFormat", true) // 新增:专门兼容旧版本Map类型的字段命名规则 .config("spark.sql.parquet.legacyMapFieldName", true) .getOrCreate();
设置这个参数后,Spark 3.5会沿用Spark 2.x的map作为Map类型的嵌套父节点,生成的Schema路径就和Spark 2.4完全一致了,完美解决Schema不匹配的问题。
2. 批量兼容现有历史数据
如果已经有大量Spark 2.4生成的Parquet文件,可以写一个Spark 3.5的批处理作业,用上述配置读取这些文件后重新写入,统一Schema格式,之后所有新的作业都用这个配置读写,就能彻底避免后续的兼容问题。
3. 验证配置是否生效
可以在代码中打印Map类型的字段路径来验证配置是否生效:
Dataset<Row> readDs = sparkSession.read().parquet("spark_all_legacy_write_enabled"); // 展开Map类型查看字段路径 readDs.select(explode(col("my_map_data")).alias("map_entry")) .select("map_entry.*") .printSchema();
如果输出的路径是my_map_data.map.key和my_map_data.map.value,就说明配置已经生效了。
注意事项
- 这个参数
spark.sql.parquet.legacyMapFieldName是Spark 3.0及以上版本新增的,专门用于兼容旧版本Parquet的Map类型字段名,在Spark 3.5中是有效的。 - 如果你后续打算彻底切换到Spark 3.x的默认格式,可以在所有作业都升级完成后,逐步移除这些Legacy参数,统一使用新的Schema格式。
内容来源于stack exchange

