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

如何在基于原生Python的AWS Glue中处理每周变更的S3文件并加载到Redshift?

我刚好处理过几乎一模一样的场景,用Glue原生Python脚本完全可以实现自动处理每周变更的S3文件并加载到Redshift,给你拆解成几个关键步骤:

1. 动态定位目标S3文件

你的文件名是按每周日期命名的(比如filenameMMDDYYYY.csv),有两种可靠的方式获取目标文件:

方式一:按日期规则生成文件名

如果能确定文件总是每周固定日期生成(比如每周日),可以直接计算对应日期来拼接文件路径:

from datetime import datetime, timedelta

# 计算最近的周日(根据你的文件日期调整逻辑,比如如果是每周一生成,就改timedelta的参数)
today = datetime.today()
last_sunday = today - timedelta(days=today.weekday() + 1)
# 格式化为MMDDYYYY的字符串
file_date_str = last_sunday.strftime("%m%d%Y")
# 拼接完整S3路径
target_s3_path = f"s3://your-bucket/your-folder/filename{file_date_str}.csv"

方式二:筛选S3前缀下的最新文件

如果日期规则可能变动,或者需要确保取到最新的文件,可以列出S3桶中符合前缀的文件,按修改时间排序取最新:

import boto3

s3_client = boto3.client('s3')
bucket_name = "your-bucket"
file_prefix = "your-folder/filename"

# 列出前缀下所有CSV文件
response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=file_prefix)
csv_files = [obj["Key"] for obj in response["Contents"] if obj["Key"].endswith(".csv")]

# 按文件最后修改时间倒序排序,取第一个(最新)
csv_files_sorted = sorted(
    csv_files,
    key=lambda x: s3_client.head_object(Bucket=bucket_name, Key=x)["LastModified"],
    reverse=True
)
target_s3_path = f"s3://{bucket_name}/{csv_files_sorted[0]}"
2. 在Glue脚本中加载S3文件

用Glue的GlueContext加载文件,支持直接解析CSV格式:

from awsglue.context import GlueContext
from pyspark.context import SparkContext

# 初始化Glue上下文
sc = SparkContext.getOrCreate()
glue_context = GlueContext(sc)

# 加载目标文件为DynamicFrame(也可以转成Spark DataFrame做更灵活的处理)
dynamic_frame = glue_context.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": [target_s3_path]},
    format="csv",
    format_options={
        "withHeader": True,  # 假设你的CSV有表头
        "separator": ",",
        "quoteChar": '"'  # 如果有带引号的字段,加上这个参数
    }
)

# 可选:转成Spark DataFrame做数据清洗/转换
spark_df = dynamic_frame.toDF()
# 比如字段类型转换:spark_df = spark_df.withColumn("amount", spark_df["amount"].cast("double"))
3. 自动调度Glue任务

要实现每周自动运行,用Glue自带的触发器或者EventBridge(原CloudWatch Events)即可:

  • Glue触发器:在Glue控制台找到你的Job,创建「定时触发器」,设置cron表达式(比如0 0 ? * MON *表示每周一凌晨0点运行,处理上周日的文件)
  • EventBridge规则:创建按cron调度的规则,目标选择你的Glue Job,同样设置每周的运行时间
4. 将数据加载到Redshift

Glue支持直接通过JDBC连接写入Redshift,有两种方式:

方式一:用Glue Catalog连接(推荐)

先在Glue控制台创建Redshift连接(存储数据库地址、用户名密码等),然后直接调用from_jdbc_conf:

glue_context.write_dynamic_frame.from_jdbc_conf(
    frame=dynamic_frame,
    catalog_connection="your-redshift-connection-name",  # 你在Glue里创建的连接名
    connection_options={
        "dbtable": "your_schema.your_target_table",
        "database": "your_redshift_db_name"
    },
    redshift_tmp_dir="s3://your-bucket/redshift-temp/"  # Redshift需要临时目录来批量加载数据
)

方式二:直接用Spark JDBC写入

如果不想用Glue Catalog连接,可以直接指定JDBC参数:

spark_df.write \
    .format("jdbc") \
    .option("url", "jdbc:redshift://your-redshift-endpoint:5439/your_db_name") \
    .option("dbtable", "your_schema.your_target_table") \
    .option("user", "your_redshift_username") \
    .option("password", "your_redshift_password") \
    .option("tempdir", "s3://your-bucket/redshift-temp/") \
    .mode("append")  # 可选:overwrite/ignore/errorifexists,根据业务需求选择
    .save()
额外注意事项
  • 权限配置:确保Glue Job的IAM角色拥有以下权限:S3读写(包括临时目录)、Redshift访问权限、Glue Catalog权限(如果用Catalog连接)
  • 错误处理:可以在脚本中加入异常捕获,比如检查文件是否存在,避免任务无意义失败:
try:
    # 检查目标文件是否存在
    s3_client.head_object(Bucket=bucket_name, Key=csv_files_sorted[0])
except s3_client.exceptions.ClientError as e:
    if e.response["Error"]["Code"] == "404":
        print("目标S3文件不存在,终止任务")
        raise Exception("Target S3 file not found")
    else:
        raise

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:48:43