如何在AWS Glue DynamicFrame中合并多Parquet文件Schema?
问题
从S3读取多个Schema不同的Parquet文件到AWS Glue DynamicFrame时,系统默认以第一个文件的Schema为准,其他文件仅保留匹配列,其余列被丢弃。如何合并所有文件的表头为统一Schema?
用户当前使用的代码:
from awsglue.context import GlueContext inputGDF = glueContext.create_dynamic_frame_from_options( connection_type = "s3", connection_options = { "paths": ["s3://bucket_name/"], "useS3ListImplementation":True, "recurse": True }, format = "parquet", transformation_ctx="inputGDF" )
方法一:启用Parquet原生Schema合并参数
Parquet格式本身支持Schema合并,直接在format_options中添加mergeSchema: True参数,即可让Glue自动合并所有文件的Schema,缺失列会填充为null:
from awsglue.context import GlueContext inputGDF = glueContext.create_dynamic_frame_from_options( connection_type = "s3", connection_options = { "paths": ["s3://bucket_name/"], "useS3ListImplementation":True, "recurse": True }, format = "parquet", format_options = {"mergeSchema": True}, # 核心参数:开启Schema合并 transformation_ctx="inputGDF" )
方法二:通过Spark DataFrame中转合并
如果第一种方法遇到类型冲突等特殊场景,可以先转为Spark DataFrame完成Schema合并,再转回Glue DynamicFrame:
from awsglue.context import GlueContext # 获取SparkSession实例 spark = glueContext.spark_session # 读取为Spark DataFrame并开启Schema合并 df = spark.read.option("mergeSchema", "true").parquet("s3://bucket_name/") # 转换为Glue DynamicFrame inputGDF = glueContext.create_dynamic_frame_from_rdd( df.rdd, schema=df.schema, transformation_ctx="inputGDF" )
注意事项
- 同名列的类型必须兼容,否则会抛出类型冲突错误。遇到这类情况可以用
ResolveChoice转换处理:from awsglue.transforms import ResolveChoice # 根据需求选择策略,比如"make_cols"将冲突列拆分为不同名称,"cast"强制转换类型 resolvedGDF = ResolveChoice.apply( frame=inputGDF, choice="make_cols", transformation_ctx="resolvedGDF" ) - 大量文件开启Schema合并会增加扫描耗时,因为需要先读取所有文件的Schema元信息。
内容的提问来源于stack exchange,提问作者user20986502
相关产品推荐
相关产品推荐

