PySpark/AWS Glue读取Parquet稀疏日期列类型推断错误问题
- 运行环境:AWS Glue 2.0,搭载Spark 2.4、Python 3
- 业务逻辑:通过
create_dynamic_frame.from_options()读取S3上存储的Parquet文件,文件源数据来自MySQL表,需要用同一套脚本迁移100+数据库,不同库的用户自定义字段(UDF)差异极大,无法硬编码指定各字段类型 - 异常现象:表中大量UDF字段填充率极低(例如某列600行仅2行有有效值),Spark会将这类列错误识别为整数类型,但该列实际是日期类型,Parquet中存储值为距Unix纪元的天数
- 当前读取代码:
dyf_full = glueContext.create_dynamic_frame.from_options(connection_type='s3', connection_options={'path': path, 'recurse': True, 'exclusions': exclusion_string}, format='parquet', transformation_ctx='df_full')
- 当前错误推断的Schema(
DATE_FIELD_2为低填充率日期列,DATE_FIELD填充率超95%可被正确识别):
root |-- CASE_ID: decimal |-- DATE_FIELD: date |-- DATE_FIELD_2: int
首先明确根因:Parquet是自带Schema的列式存储格式,Spark读取Parquet默认不会扫描全量数据做类型推断,而是直接读取文件头写入的Schema。你遇到的类型错配,本质是上游写入Parquet文件时,因为对应列非空值占比太低,写入端就把该列的Schema标记成了整数类型,和读取端的采样推断没有关系。
以下是不需要硬编码字段、可适配多库迁移场景的可行方案:
方案1:开启Parquet Schema合并,强制拉取所有文件的Schema做对齐
读取时在format_options中开启Schema合并配置,Spark会遍历路径下所有Parquet文件的元数据,合并所有分片的字段类型,只要任意一个文件分片中该列存在非空的日期类型值,合并后就会得到正确的日期类型,完全兼容DynamicFrame:dyf_full = glueContext.create_dynamic_frame.from_options(connection_type='s3', connection_options={'path': path, 'recurse': True, 'exclusions': exclusion_string}, format='parquet', format_options={ "mergeSchema": "true" }, transformation_ctx='df_full')Glue 2.0内置的Spark 2.4默认支持该参数,不需要额外修改Spark配置。
方案2:通用字段类型自动修正,兜底Schema合并未覆盖的场景
如果开启Schema合并后仍有个别低填充率列被识别为整数,可以加一层通用转换逻辑:遍历所有字段,自动识别符合日期存储特征的整数字段,统一转换为日期类型,不需要提前知道字段名:from pyspark.sql import functions as F from pyspark.sql.types import IntegerType, DateType from awsglue.dynamicframe import DynamicFrame df = dyf_full.toDF() for field in df.schema.fields: # 仅处理整数类型字段 if isinstance(field.dataType, IntegerType): # 统计该列非空值的范围,判断是否符合日期天数的取值区间(1970-01-01到2100-01-01对应0~47482天) stat = df.agg( F.max(field.name).alias("max_v"), F.min(field.name).alias("min_v"), F.count(field.name).alias("nn_cnt") ).first() if stat.nn_cnt > 0 and stat.min_v >= 0 and stat.max_v <= 47482: # 把天数偏移转为标准日期类型 df = df.withColumn( field.name, F.date_add(F.lit("1970-01-01").cast(DateType()), F.col(field.name)) ) # 转换回DynamicFrame供后续Glue逻辑使用 dyf_full = DynamicFrame.fromDF(df, glueContext, "fixed_dyf")这段逻辑对所有库通用,不管UDF字段名是什么,只要是用Unix天数存储的日期字段,都会被自动识别转换,不会误转普通整数字段(普通业务整数字段一般会超过2100年对应的天数值,不会落入判定区间)。
方案3:上游写入环节修正Schema(如果有权限调整写入作业)
如果可以修改生成Parquet文件的上游作业,直接在写入前对低填充率列做一次全量类型校验,确保写入Parquet的Schema和实际类型一致,从根源避免类型错配问题。
内容的提问来源于stack exchange,提问作者whit.the.engineer

