AWS Glue通过create_dynamic_frame.from_options获取S3文件夹分区的方法
解决方案
方案1:使用Glue原生分区识别
如果你的S3分区层级是固定顺序的,不管是否是分区键=分区值的Hive风格命名,都可以直接在create_dynamic_frame.from_options的参数中配置识别路径分区,不需要数据本身包含分区字段,也不需要依赖爬网程序:
connection_type指定为s3connection_options中传入数据的根S3路径,开启递归读取,同时按路径从上到下的层级顺序在partitionKeys中指定你要给每个分区层级设置的字段名
示例代码:
from awsglue.context import GlueContext from pyspark.context import SparkContext glueContext = GlueContext(SparkContext.getOrCreate()) dynamic_frame = glueContext.create_dynamic_frame.from_options( connection_type="s3", connection_options={ "paths": ["s3://你的存储桶/数据根路径/"], "recurse": True, # 示例路径结构为s3://bucket/date=20240101/hour=12/,则按层级顺序填分区键名 "partitionKeys": ["date", "hour"] }, format="parquet", # 替换为你的实际数据格式,比如json、csv等 format_options={} # 按需补充格式参数,比如压缩格式、分隔符等 )
生成的动态框架会自动将路径对应层级的文件夹名作为分区字段的值,添加到数据schema中。
方案2:手动解析S3路径提取分区(适配任意不规则路径)
如果你的分区路径结构不固定、或者有自定义命名规则,可以通过input_file_name()函数获取每条数据对应的源文件S3全路径,自行拆分提取分区字段,不受任何路径格式限制:
示例代码:
from pyspark.sql.functions import input_file_name, split, element_at # 先把动态框架转为Spark DataFrame方便做字段处理 df = dynamic_frame.toDF() # 新增文件路径字段,按实际路径层级拆分提取分区字段 df_with_partition = df.withColumn("file_full_path", input_file_name()) \ .withColumn("path_parts", split("file_full_path", "/")) \ # 示例:假设路径倒数第3层是年、倒数第2层是月、倒数第1层是日,按自己的路径结构调整下标即可 .withColumn("year", element_at("path_parts", -3)) \ .withColumn("month", element_at("path_parts", -2)) \ .withColumn("day", element_at("path_parts", -1)) \ .drop("file_full_path", "path_parts") # 按需删除中间临时字段 # 如果后续需要用Glue API处理,再转回动态框架即可 from awsglue.dynamicframe import DynamicFrame final_dynamic_frame = DynamicFrame.fromDF(df_with_partition, glueContext, "final_dynamic_frame")
内容的提问来源于stack exchange,提问作者123
相关产品推荐
相关产品推荐

