如何调整AWS Glue Jobs输出至S3的分区目录结构?
问题
我正在运行AWS Glue Jobs处理多个关联表,这些表使用ts(timestamp)作为分区键。默认情况下,每个Glue Job在S3中写入的输出文件目录结构如下(以指定表和时间戳为例):
s3://someBucket/someFolder/table1/ts=2023-03-08T21:20:17Z/data*.parquet s3://someBucket/someFolder/2023-03-08T21:20:17Z/table2/data*.parquet s3://someBucket/someFolder/2023-03-08T21:20:17Z/table3/data*.parquet
由于所有表共享相同的时间戳,我希望改为如下目录结构:
s3://someBucket/someFolder/2023-03-08T21:20:17Z/table1/data*.parquet s3://someBucket/someFolder/2023-03-08T21:20:17Z/table2/data*.parquet s3://someBucket/someFolder/2023-03-08T21:20:17Z/table3/data*.parquet
请问该需求是否可行?如果可行,具体该如何实现?
补充代码
以下是相关Glue Job代码(与AWS生成的简单Python Glue Jobs代码差异不大):
# Script generated for node S3 bucket S3bucket_node1 = glueContext.create_dynamic_frame.from_options( connection_type="s3", format="json", connection_options={ "paths": [s3Source], }, format_options={ "multiline": False }, transformation_ctx="S3bucket_node1", ) # Script generated for node S3 bucket S3bucket_node3 = glueContext.getSink( path=s3Destination, connection_type="s3", updateBehavior="UPDATE_IN_DATABASE", partitionKeys=["ts"], compression="snappy", enableUpdateCatalog=True, transformation_ctx="S3bucket_node3", ) S3bucket_node3.setCatalogInfo( catalogDatabase="some_data_catalog", catalogTableName="table1" ) S3bucket_node3.setFormat("glueparquet") S3bucket_node3.writeFrame(S3bucket_node1) job.commit()
我找不到修改S3分区格式的方法,很困扰。
可行方案
这个需求完全可以实现,核心思路是放弃Glue自动生成的partitionKeys分区逻辑,手动构造输出路径并写入数据,同时确保Glue Data Catalog的元数据能正确关联到新的路径结构。具体步骤如下:
1. 移除自动分区配置
在getSink中删除partitionKeys=["ts"]参数,因为我们不再依赖Glue自动按分区键生成目录。
2. 手动构造输出路径
将时间戳(ts值)直接拼接到输出路径中,再加上表名。比如原来的s3Destination如果是s3://someBucket/someFolder/table1/,现在要改成s3://someBucket/someFolder/{ts_value}/table1/。
注意:你需要提前获取本次Job处理的ts值——可以从输入数据中提取,或者作为Job参数传入(更可靠,避免数据中ts值不一致)。
3. 调整Sink配置并写入
修改后的代码示例(以table1为例):
# 假设已经获取到当前处理的时间戳值,比如从Job参数传入 current_ts = "2023-03-08T21:20:17Z" # 构造新的输出路径 new_s3_destination = f"s3://someBucket/someFolder/{current_ts}/table1/" # 创建Sink,移除partitionKeys参数 S3bucket_node3 = glueContext.getSink( path=new_s3_destination, connection_type="s3", updateBehavior="UPDATE_IN_DATABASE", compression="snappy", enableUpdateCatalog=True, transformation_ctx="S3bucket_node3", ) S3bucket_node3.setCatalogInfo( catalogDatabase="some_data_catalog", catalogTableName="table1" ) S3bucket_node3.setFormat("glueparquet") # 写入时确保数据中保留ts字段,以便Catalog识别分区 S3bucket_node3.writeFrame(S3bucket_node1) job.commit()
4. 维护Glue Data Catalog分区
因为手动构造了路径,需要确保Glue Catalog中的表分区能正确映射到新路径,有两种方式:
- 手动添加分区:执行SQL
ALTER TABLE table1 ADD PARTITION (ts='2023-03-08T21:20:17Z') LOCATION 's3://someBucket/someFolder/2023-03-08T21:20:17Z/table1/',可通过Glue Job调用boto3操作Athena执行这条语句。 - 自动发现分区:在Glue表设置中开启分区自动发现,或定期运行Glue Crawler扫描新路径,让Crawler自动识别并添加分区。
关键注意事项
- 所有关联表必须使用相同的ts值构造路径,才能实现按时间戳聚合目录的效果。
- 数据中必须保留
ts字段,否则Glue Catalog无法将路径与分区键关联。 - 如果批量处理多个时间戳,需循环处理每个ts值,构造对应路径并写入数据。
内容的提问来源于stack exchange,提问作者Alvaro Mendez
相关产品推荐
相关产品推荐

