Glue PySpark写入Parquet报PlainLongDictionary类型错误
问题描述
累计排查2天未定位根因,AWS支持团队未提供有效协助,问题场景如下:
- Glue中
provisioned库下定义了目标表provisioned.customer
- 上游依赖两张表
curated.customer_consents、curated.customer,所有字段均已设置为string类型,表结构如下:

- 作业运行环境为Glue 3.0,PySpark作业代码如下:
# more imports... import pyspark.sql.types as T from awsglue.dynamicframe import DynamicFrame from pyspark.sql import functions as F def fudf(val): return functools.reduce(lambda x, y: x+y, val) flattenUdf = F.udf(fudf, T.StringType()) customer_df = spark.read.table('curated.customer') core_customer_consents_df = spark.read.table('curated.core_customer_consents') core_customer_consents_df = core_customer_consents_df.groupBy("tscid").agg(F.collect_list('purpose').alias('__purposes__')) merged = customer_df.join(core_customer_consents_df, ['tscid'], how='inner') merged = merged.select("*", flattenUdf("__purposes__").alias("cirrus_purposes")) merged = merged.select([c for c in merged.columns if c not in {'__purposes__'}]) logger.info(merged) glueContext.write_dynamic_frame.from_catalog( frame=DynamicFrame.fromDF(merged, glueContext, "a"), database = 'provisioned', table_name = 'customer', transformation_ctx = "datasource0", )
- 作业日志打印的DataFrame schema如下,确认所有字段均为string类型,与目标表字段顺序、类型完全匹配:
2022-06-15 08:35:26 INFO logger: DataFrame[tscid: string, first_name: string, last_name: string, gender: string, birthdate: string, qualification: string, country: string, estimated_annual_earning: string, cirrus_purposes: string]
- 写入阶段抛出如下报错:
An error occurred while calling o149.pyWriteDynamicFrame. org.apache.parquet.column.values.dictionary.PlainValuesDictionary$PlainLongDictionary

根因分析
这个报错和你反复校验的DataFrame schema、Glue Catalog登记的表schema没有任何关系,是Glue写入Parquet的默认机制导致的坑:write_dynamic_frame.from_catalog写入Parquet格式表时,默认会拉取目标表对应S3路径下所有已存在的历史Parquet文件的footer元数据做自动schema合并,不会完全采信传入的DF schema或者Catalog中登记的schema。
你的目标表S3路径下残留了之前写入的老Parquet文件,这些文件中至少有一个字段是Long类型(大概率是早期建表时把tscid或其他字段设为bigint类型时写入的文件)。Parquet的列字典编码和字段类型强绑定,Long类型生成的PlainLongDictionary无法被识别为String类型,直接抛出类型转换异常。
解决方案
按优先级从高到低选择对应方案即可:
- 方案1(最快修复,适合允许清空历史数据的场景):直接清空
provisioned.customer表对应S3路径下的所有历史文件,再重新运行作业。清空后不存在旧schema的Parquet文件干扰,作业会正常写入全String类型的新Parquet文件,后续写入也不会再触发该异常。 - 方案2(适合需要保留历史数据的场景):写入时关闭Parquet自动合并schema的配置,强制使用当前传入的DataFrame schema写入,修改写入代码如下:
glueContext.write_dynamic_frame.from_catalog( frame=DynamicFrame.fromDF(merged, glueContext, "a"), database = 'provisioned', table_name = 'customer', transformation_ctx = "datasource0", additional_options = { "mergeSchema": "false" } )
同时可以在作业的Spark配置项中添加spark.sql.parquet.mergeSchema=false,双重保险关闭schema自动合并逻辑。
- 性能优化(非强制,可解决自定义UDF的稳定性问题):替换自定义Python UDF为Spark内置
concat_ws函数实现数组转字符串,性能提升10~100倍,且避免空值处理异常,修改对应逻辑如下:
# 删掉自定义fudf、flattenUdf的定义,替换聚合逻辑 core_customer_consents_df = core_customer_consents_df.groupBy("tscid").agg( F.concat_ws('', F.collect_list('purpose')).alias('cirrus_purposes') ) # 后续不需要再做drop __purposes__列、UDF处理步骤,直接join即可 merged = customer_df.join(core_customer_consents_df, ['tscid'], how='inner')
内容的提问来源于stack exchange,提问作者Tomás Denis Reyes Sánchez
相关产品推荐
相关产品推荐

