如何在AWS Glue可视化ETL中提取CSV文件名并处理单文件校验错误
解决方案:AWS Glue ETL 关联文件元数据与错误处理
1. 从DynamicFrame提取每行对应的源文件名
可视化ETL方式
在Glue可视化ETL的Glue Catalog数据源节点中,找到Advanced properties下的Add source file name to output选项并勾选,系统会自动在输出的DynamicFrame中添加source_file字段,记录每行数据对应的源文件完整S3路径。
脚本ETL方式
读取数据时通过配置参数保留元数据,再用Spark函数提取文件名:
from awsglue.context import GlueContext from pyspark.sql.functions import input_file_name glueContext = GlueContext(spark.sparkContext) # 从Catalog读取数据,启用递归扫描 dynamic_frame = glueContext.create_dynamic_frame.from_catalog( database="your_db_name", table_name="your_table_name", additionalOptions={"recurse": True} ) # 转换为DataFrame并添加文件名列 df_with_file = dynamic_frame.toDF().withColumn("source_file", input_file_name()) # 转回DynamicFrame以适配后续可视化节点操作 enhanced_dynamic_frame = glueContext.create_dynamic_frame.from_df( df_with_file, glueContext, "enhanced_frame_with_filename" )
2. 在Glue可视化ETL中按文件执行校验
可视化ETL无法直接按文件分组校验,但可通过以下步骤实现:
- 先按上述方法给数据集添加
source_file字段 - 添加自定义转换节点,编写PySpark代码按文件分组执行校验逻辑:
from pyspark.sql.functions import lit def transform(glueContext, dfc): df = dfc.select(list(dfc.keys())[0]) # 定义单文件校验逻辑 def validate_single_file(group_df): file_path = group_df.select("source_file").first()[0] # 示例:检查关键字段是否存在空值 critical_null_count = group_df.filter("user_id IS NULL OR order_date IS NULL").count() if critical_null_count > 0: return group_df.withColumn("validation_status", lit("FAILED")) else: return group_df.withColumn("validation_status", lit("PASSED")) # 按文件分组执行校验 validated_df = df.groupBy("source_file").flatMapGroups(validate_single_file) return [glueContext.create_dynamic_frame.from_df(validated_df, glueContext, "validated_frame")]
- 添加拆分节点,根据
validation_status字段将数据拆分为“校验通过”和“校验失败”两个分支。
3. 编程方式隔离并移动错误文件
步骤1:提取错误文件的唯一路径
从校验失败的DynamicFrame中提取不重复的源文件路径:
failed_dynamic_frame = ... # 拆分节点输出的失败数据集 failed_files_df = failed_dynamic_frame.toDF().select("source_file").distinct() failed_file_paths = [row.source_file for row in failed_files_df.collect()]
步骤2:移动错误文件到S3失败文件夹
使用Boto3实现文件移动,需确保Glue作业的IAM角色拥有S3的CopyObject和DeleteObject权限:
import boto3 s3_client = boto3.client("s3") source_bucket = "your_source_bucket_name" failure_folder_prefix = "data-failure/" for file_path in failed_file_paths: # 解析S3文件键(去除s3://bucket/前缀) file_key = file_path.replace(f"s3://{source_bucket}/", "") # 构建目标路径 dest_key = f"{failure_folder_prefix}{file_key.split('/')[-1]}" # 复制文件到失败文件夹 s3_client.copy_object( Bucket=source_bucket, CopySource={"Bucket": source_bucket, "Key": file_key}, Key=dest_key ) # 删除原文件(可选,根据业务需求决定) s3_client.delete_object(Bucket=source_bucket, Key=file_key)
内容的提问来源于stack exchange,提问作者coderman
相关产品推荐
相关产品推荐

