如何在AWS Glue动态框架中添加分区dt字段(不使用目录数据库)
为AWS Glue DynamicFrame添加分区字段dt(无需目录数据库)
你可以通过以下方式在不依赖Glue目录数据库的前提下,从文件路径中提取分区字段dt并添加到DynamicFrame中:
实现步骤与代码示例
from pyspark.sql.functions import input_file_name, regexp_extract from awsglue.dynamicframe import DynamicFrame # 原数据加载代码 dynamic_frame = glueContext.create_dynamic_frame.from_options( connection_type="s3", format='parquet', connection_options={ "paths": ["s3://my_bucket/root/"], "recurse": True, }, transformation_ctx="S3bucket_node1", ) # 1. 将DynamicFrame转换为Spark DataFrame,方便路径解析操作 df = dynamic_frame.toDF() # 2. 从文件路径提取dt字段:用正则匹配dt=后的日期内容 df_with_dt = df.withColumn( "dt", regexp_extract(input_file_name(), r'dt=([^/]+)', 1) ) # 3. 转回DynamicFrame,保持Glue ETL操作兼容性 dynamic_frame_with_dt = DynamicFrame.fromDF(df_with_dt, glueContext, "dynamic_frame_with_dt")
代码说明
input_file_name():获取每条数据对应的源S3文件路径,例如s3://my_bucket/root/dt=2022-24-12/file.parquetregexp_extract():通过正则表达式dt=([^/]+)精准提取dt=后到下一个/前的日期值,适配你的分区路径格式- 若日期格式固定,也可改用
split()函数替代正则,示例如下:from pyspark.sql.functions import split df_with_dt = df.withColumn( "dt", split(split(input_file_name(), "dt=")[1], "/")[0] )
处理完成后,dynamic_frame_with_dt会包含原有的column_A、Column_B字段,以及对应分区的dt字段。
内容的提问来源于stack exchange,提问作者Shlomi Schwartz
相关产品推荐
相关产品推荐

