AWS Glue技术问询:如何为输出添加源文件名列?
在Glue Python作业中添加原始文件名列的方法
我刚好处理过类似的需求,给你分享几个实用的方案,都是基于Glue Python作业的特性(毕竟Glue底层是Spark,很多Spark的函数都能直接用):
方法一:利用Spark内置函数获取完整文件路径
这是最直接的方式,Spark提供了input_file_name()函数,可以直接获取每条数据对应的源文件完整S3路径。步骤如下:
- 先通过Glue Catalog读取数据到DynamicFrame(也就是你爬虫生成的表)
- 把DynamicFrame转换成Spark DataFrame(方便使用Spark的函数)
- 添加包含源文件名的新列
- 可选:转回DynamicFrame继续使用Glue的API(如果后续需要)
示例代码:
from pyspark.sql.functions import input_file_name from awsglue.dynamicframe import DynamicFrame # 从Glue Catalog读取爬虫生成的表 source_dyf = glueContext.create_dynamic_frame.from_catalog( database="你的数据库名", table_name="你的表名" ) # 转换为Spark DataFrame source_df = source_dyf.toDF() # 添加原始文件路径列(完整S3路径) df_with_file = source_df.withColumn("original_file_path", input_file_name()) # 可选:转回DynamicFrame(如果后续要用Glue的写操作或其他API) dyf_with_file = DynamicFrame.fromDF(df_with_file, glueContext, "dyf_with_file")
方法二:提取纯文件名(去掉S3路径)
如果只需要文件名而不是完整路径,可以用regexp_extract函数配合input_file_name()来提取:
from pyspark.sql.functions import regexp_extract, input_file_name # 在上一步的source_df基础上添加纯文件名列 df_with_filename = source_df.withColumn( "original_filename", regexp_extract(input_file_name(), ".*/(.*)$", 1) )
这里的正则表达式.*/(.*)$会匹配最后一个/之后的所有内容,也就是纯文件名。
方法三:直接用Spark读取S3文件(不通过Glue Catalog)
如果你的场景不需要依赖Glue Catalog,直接读取S3文件的话,同样可以用上面的函数:
# 直接读取S3上的CSV文件 source_df = spark.read.format("csv") \ .option("header", "true") \ .load("s3://你的源Bucket路径/") # 添加文件名列 df_with_filename = source_df.withColumn("original_filename", input_file_name())
补充说明
你提到想通过awsglue.job包获取元数据,其实这个包主要是用来获取作业本身的元数据(比如作业ID、运行参数、执行时间等),而源文件名属于数据层面的元数据,所以需要借助Spark的内置函数来获取,这也是官方推荐的方式。
最后把处理好的数据写回S3即可,比如:
# 写回CSV到目标S3桶 glueContext.write_dynamic_frame.from_options( frame=dyf_with_file, connection_type="s3", connection_options={"path": "s3://你的目标Bucket路径/"}, format="csv", format_options={"writeHeader": True} )
内容的提问来源于stack exchange,提问作者markwatson
相关产品推荐
相关产品推荐

