You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 23:12:35