You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Glue PySpark写入Parquet报PlainLongDictionary类型错误

问题描述

累计排查2天未定位根因,AWS支持团队未提供有效协助,问题场景如下:

  • Glue中provisioned库下定义了目标表provisioned.customer
    provisioned.customer表配置截图
  • 上游依赖两张表curated.customer_consents、curated.customer,所有字段均已设置为string类型,表结构如下:
    curated.customer_consents表结构截图
    curated.customer表结构截图
  • 作业运行环境为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 23:16:02