AWS Glue作业分区写S3报AttributeError: '_jdf'属性不存在怎么解决
AWS Glue分区写入S3报错解决及实现方案
报错根因
df.write.partitionBy("last_modified_date").mode("overwrite").parquet("s3://***") 是Spark的写入动作,执行完成后返回的是None类型,你后续尝试将这个None值转换为DynamicFrame,就会触发AttributeError: 'NoneType' object has no attribute '_jdf'报错。如果不需要对写入后的结果做二次处理,最后一行转换代码完全多余,可以直接删除。
实现方案
方案1:修正后的PySpark原生写法
如果不需要用到Glue特有的动态帧能力,直接用Spark DataFrame写入即可:
from pyspark.sql.functions import current_timestamp # 从数据目录读取表转DataFrame datasource0 = glueContext.create_dynamic_frame.from_catalog(database='***', table_name= d['source_table']) df = datasource0.toDF().withColumn('logged_at_utc', current_timestamp()) # 直接按指定字段分区写入S3,不需要接收返回值 df.write.partitionBy("last_modified_date").mode("overwrite").parquet("s3://***")
如果写入后确实需要将落地的数据重新加载为动态帧做后续处理,需要重新读取S3路径的文件,不能直接使用write操作的返回值转换。
方案2:Glue DynamicFrame原生写入(更推荐)
使用Glue自带的Sink接口写入,适配Glue生态,支持自动同步分区元数据到数据目录,无需手动执行MSCK命令修复分区:
from pyspark.sql.functions import current_timestamp from awsglue.dynamicframe import DynamicFrame # 读取源数据并添加字段 datasource0 = glueContext.create_dynamic_frame.from_catalog(database='***', table_name= d['source_table']) df = datasource0.toDF().withColumn('logged_at_utc', current_timestamp()) dyf = DynamicFrame.fromDF(df, glueContext, "processed_dyf") # 配置S3写入参数 s3_sink = glueContext.getSink( path="s3://你的目标存储路径", connection_type="s3", updateBehavior="UPDATE_IN_DATABASE", partitionKeys=["last_modified_date"], enableUpdateCatalog=True, catalogDatabase="你的目标数据库名", # 可选,需要同步元数据时填写 catalogTableName="你的目标表名", # 可选,需要同步元数据时填写 ) s3_sink.setFormat("parquet") s3_sink.setMode("overwrite") # 执行写入 s3_sink.writeFrame(dyf)
该方案写入完成后会自动在Glue数据目录创建/更新目标表的元数据和分区信息,后续可直接通过数据目录查询写入后的分区表。
内容的提问来源于stack exchange,提问作者averma
相关产品推荐
相关产品推荐

