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

AWS Glue作业路径错误:S3文件读取并加载Redshift表问题

AWS Glue作业路径错误修复及代码优化

核心路径错误分析

  • 文件名列表获取错误:原代码中file_list=list(file_names)将DataFrame转为列对象列表,而非实际的S3文件路径字符串列表,导致后续循环无法读取正确文件。
  • 循环读取文件时路径硬编码:paths: ["obj"]是使用字符串"obj"作为路径,而非遍历的变量obj的实际值,导致S3路径无效。

修正后的完整代码

import sys
import logging
from datetime import datetime,timedelta,date
from pyspark.sql.functions import input_file_name
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.dynamicframe import DynamicFrame
from awsglue.job import Job
import boto3
import pandas as pd
import json
from botocore.exceptions import ClientError

if __name__ == "__main__":
   
    args = getResolvedOptions(sys.argv, ['TempDir','JOB_NAME','secret_name'])
    secret_name = args['secret_name']
    temp_dir = args['TempDir']
   
    sc = SparkContext()
    glueContext = GlueContext(sc)
    spark = glueContext.spark_session
    job = Job(glueContext)
    job.init(args['JOB_NAME'], args)
   
    # 获取Redshift凭证
    session = boto3.session.Session()
    client = session.client(
        service_name='secretsmanager',
        region_name='us-east-1'
    )
    
    secret = None
    try:
        get_secret_value_response = client.get_secret_value(
            SecretId=secret_name
        )
        if 'SecretString' in get_secret_value_response:
            secret = json.loads(get_secret_value_response['SecretString'])
    except ClientError as e:
        print("获取凭证错误:",e)
        # 凭证获取失败时直接终止作业
        job.commit()
        sys.exit(1)
           
    # 提取用户名和密码
    redshiftUserID = secret.get('username')
    redshiftPassword = secret.get('password')
   
    print("作业启动时间: " + str(datetime.now()))
   
    # JDBC连接URL
    jdbc_url = "jdbc:redshift://rcm-analytics-dev2.cvxngrljq9pn.us-east-1.redshift.amazonaws.com:5439/rcmanalyticsdev2"
   
    # 读取S3目录下所有CSV文件(仅用于提取文件名,不加载全量数据)
    Client_files = glueContext.create_dynamic_frame.from_options(
        connection_type="s3",
        format="csv",
        connection_options={
            "paths": ["s3://rcma-chc-dev-analytics-cloud/merge/input"],
            'recurse':True
        },
        format_options={'withHeader': True},
        transformation_ctx = "Client_files"
    )

    print("S3文件读取完成: " + str(datetime.now()))
    
    # 提取所有唯一文件名
    file_names_df = Client_files.toDF().withColumn("Filename", input_file_name()).select("Filename").distinct()
    
    # 打印文件名列表
    print("所有文件路径:")
    file_names_df.show(truncate=False)
    
    # 转换为实际的S3路径字符串列表
    file_list = file_names_df.rdd.map(lambda row: row.Filename).collect()
    
    # 打印每个文件路径
    for file_path in file_list:
        print(file_path)

    # 提前读取Redshift目标表Schema,避免循环内重复读取
    conn_options = {
        "url": jdbc_url,
        "database":"rcmanalyticsdev2",
        "dbtable": "ODS_EPREMIS.CSV_MERGE_INTERQUAL",
        "user": redshiftUserID,
        "password": redshiftPassword,
        "redshiftTmpDir": args["TempDir"]
    }
    target_frame = glueContext.create_dynamic_frame_from_options("redshift", conn_options)
    target_column_names = [field.name for field in target_frame.schema().fields]
    target_column_count = len(target_column_names)

    # 遍历每个文件执行ETL
    for obj in file_list:
        print(f"开始处理文件: {obj}")
        # 读取单个S3文件
        each_clientfiles = glueContext.create_dynamic_frame.from_options(
            connection_type="s3",
            format="csv",
            connection_options={
                "paths": [obj],  # 使用变量obj的实际值,而非字符串"obj"
                'recurse':False  # 单个文件不需要递归
            },
            format_options={'withHeader': True},
            transformation_ctx = f"each_clientfiles_{obj.split('/')[-1]}"
        )
        
        # 转换为DataFrame
        frame = each_clientfiles.toDF()
        source_column_count = len(frame.columns)
        
        try:
            # 校验列数是否匹配
            if source_column_count != target_column_count:
                raise ValueError(f"文件{obj}列数不匹配,源文件列数:{source_column_count},目标表列数:{target_column_count}")
            
            # 对齐列名(确保列顺序和名称与目标表一致)
            frame = frame.toDF(*target_column_names)
            final_dynamic_frame = DynamicFrame.fromDF(frame, glueContext, "final")

            # 将数据写入Redshift
            redshift_write = glueContext.write_dynamic_frame.from_jdbc_conf(
                frame=final_dynamic_frame,
                catalog_connection = "acuity-redshift",
                redshift_tmp_dir = args["TempDir"],  # 使用作业参数中的临时目录
                connection_options=conn_options,
                transformation_ctx = f"redshift_write_{obj.split('/')[-1]}"
            )
            print(f"文件{obj}写入Redshift成功")
        except Exception as e:
            print(f"处理文件{obj}时出错:", e)
        finally:    
            print(f"文件{obj}处理完成: {str(datetime.now())}")
    
    print("Glue作业全部完成: ", str(datetime.now()))
    job.commit()

其他关键优化点

  • 凭证异常处理增强:添加凭证获取失败时的终止逻辑,避免后续因secret未定义导致报错。
  • 目标表Schema提前读取:将Redshift目标表Schema读取移到循环外,避免重复调用Redshift,提升作业效率。
  • 列数显式校验:新增列数匹配检查,提前抛出明确异常,替代原代码中toDF隐式报错。
  • 临时目录使用作业参数:替换硬编码的临时目录为作业传入的TempDir参数,符合AWS Glue最佳实践。
  • 单个文件读取关闭递归:读取单个文件时设置recurse:False,避免不必要的目录遍历。

内容的提问来源于stack exchange,提问作者Sumit Gupta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 20:44:56