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

AWS Glue动态路径S3 CSV同步至Redshift:Catalog配置与脚本实现

解决AWS Glue同步S3动态日期路径到Redshift的问题

我刚好处理过类似的动态分区同步场景,下面分两部分帮你搞定这个需求:


一、在Glue数据目录中添加动态S3路径(仅最新日期文件夹)

你的S3路径是按d=YYYY-MM-DD分区的,我们可以利用Glue的分区表特性来管理,同时确保只同步最新的分区:

方法1:手动创建分区表 + 动态注册最新分区

  1. 先创建基础Glue表

    • 登录Glue控制台,进入数据目录 → 表 → 添加表
    • 数据源选S3,设置表的根路径为s3://data-dl/abc/
    • 定义分区键:添加名为d的分区列,类型选字符串(格式为YYYY-MM-DD)
    • 配置CSV格式:设置分隔符、是否包含表头(如果你的CSV有表头,记得勾选“跳过表头行”),完成表创建。
  2. 动态注册最新分区到Glue表
    在ETL作业运行前,先通过代码获取S3上最新的日期分区,然后将其注册到Glue表中(如果不存在的话)。可以用boto3实现:

    import boto3
    from datetime import datetime
    
    def register_latest_partition():
        # 初始化S3和Glue客户端
        s3_client = boto3.client('s3')
        glue_client = boto3.client('glue')
        
        bucket = "data-dl"
        prefix = "abc/"
        db_name = "your_glue_database"
        table_name = "your_glue_table"
    
        # 列出所有d=前缀的分区文件夹
        response = s3_client.list_objects_v2(Bucket=bucket, Prefix=prefix, Delimiter='/')
        partitions = []
        for prefix_obj in response.get('CommonPrefixes', []):
            folder = prefix_obj['Prefix'].strip('/').split('/')[-1]
            if folder.startswith('d='):
                date_str = folder.split('=')[1]
                # 验证日期格式有效性
                try:
                    datetime.strptime(date_str, '%Y-%m-%d')
                    partitions.append(date_str)
                except ValueError:
                    continue
    
        if not partitions:
            raise Exception("No valid date partitions found in S3")
    
        # 获取最新日期
        latest_date = max(partitions)
        partition_location = f"s3://{bucket}/{prefix}d={latest_date}/"
    
        # 检查分区是否已存在,不存在则创建
        try:
            glue_client.get_partition(
                DatabaseName=db_name,
                TableName=table_name,
                PartitionValues=[latest_date]
            )
            print(f"Partition d={latest_date} already exists in Glue Data Catalog")
        except glue_client.exceptions.EntityNotFoundException:
            # 构建分区的存储描述符(需和你表的结构匹配)
            glue_client.create_partition(
                DatabaseName=db_name,
                TableName=table_name,
                PartitionInput={
                    "Values": [latest_date],
                    "StorageDescriptor": {
                        "Location": partition_location,
                        "InputFormat": "org.apache.hadoop.mapred.TextInputFormat",
                        "OutputFormat": "org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat",
                        "SerdeInfo": {
                            "SerializationLibrary": "org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe",
                            "Parameters": {
                                "field.delim": ",",
                                "skip.header.line.count": "1"  # 如果CSV有表头请保留
                            }
                        },
                        # 这里填写你的表字段,比如:
                        "Columns": [
                            {"Name": "col1", "Type": "string"},
                            {"Name": "col2", "Type": "int"},
                            # 其他字段...
                        ]
                    }
                }
            )
            print(f"Successfully registered partition d={latest_date}")
    

方法2:配置Glue爬虫自动爬取最新分区

如果你不想手动写注册代码,可以用Glue爬虫:

  • 创建爬虫,数据源设为s3://data-dl/abc/,目标数据库和表选你创建的分区表
  • 在爬虫配置中,开启增量爬取,设置爬取频率为每日(和你的文件生成时间匹配)
  • 爬虫会自动检测新增的d=YYYY-MM-DD文件夹,并将其添加为Glue表的分区

二、在Glue PySpark脚本中声明动态路径

有两种常见方式,根据你的需求选择:

方式1:直接读取最新S3路径(不依赖Glue数据目录)

如果不需要通过Glue表读取,直接构造最新日期的S3路径读取CSV:

import boto3
from datetime import datetime

def get_latest_s3_path():
    s3_client = boto3.client('s3')
    bucket = "data-dl"
    prefix = "abc/"
    
    response = s3_client.list_objects_v2(Bucket=bucket, Prefix=prefix, Delimiter='/')
    partitions = []
    for prefix_obj in response.get('CommonPrefixes', []):
        folder = prefix_obj['Prefix'].strip('/').split('/')[-1]
        if folder.startswith('d='):
            date_str = folder.split('=')[1]
            try:
                datetime.strptime(date_str, '%Y-%m-%d')
                partitions.append(date_str)
            except ValueError:
                continue
    
    if not partitions:
        raise Exception("No valid date partitions found")
    
    latest_date = max(partitions)
    return f"s3://{bucket}/{prefix}d={latest_date}/"

# 获取最新路径并读取数据
latest_path = get_latest_s3_path()
df = spark.read.csv(
    latest_path,
    header=True,       # 如果CSV有表头设为True
    inferSchema=False, # 建议手动指定Schema,避免自动推断出错
    sep=',',
    schema=your_custom_schema  # 替换为你的表结构,比如 StructType([StructField("col1", StringType()), ...])
)

# 后续同步到Redshift的代码
df.write \
    .format("com.databricks.spark.redshift") \
    .option("url", "jdbc:redshift://your-redshift-cluster:5439/your-db?user=your-user&password=your-password") \
    .option("dbtable", "target_redshift_table") \
    .option("tempdir", "s3://your-temp-bucket/temp-folder/")  # Redshift需要临时存储路径
    .mode("append")  # 根据需求选append、overwrite等
    .save()

方式2:通过Glue分区表读取最新分区

如果你已经将最新分区注册到Glue表,可以直接过滤读取:

import boto3
from datetime import datetime

def get_latest_date():
    s3_client = boto3.client('s3')
    bucket = "data-dl"
    prefix = "abc/"
    
    response = s3_client.list_objects_v2(Bucket=bucket, Prefix=prefix, Delimiter='/')
    partitions = []
    for prefix_obj in response.get('CommonPrefixes', []):
        folder = prefix_obj['Prefix'].strip('/').split('/')[-1]
        if folder.startswith('d='):
            date_str = folder.split('=')[1]
            try:
                datetime.strptime(date_str, '%Y-%m-%d')
                partitions.append(date_str)
            except ValueError:
                continue
    
    return max(partitions) if partitions else None

# 获取最新日期并读取Glue表的对应分区
latest_date = get_latest_date()
df = spark.read.table("your_glue_database.your_glue_table").filter(f"d = '{latest_date}'")

# 后续同步到Redshift的代码和上面一致

注意事项

  • 确保Glue作业的IAM角色有S3读取权限、Glue数据目录操作权限,以及Redshift的写入权限
  • YYYY-MM-DD格式的字符串可以直接用max()获取最新日期,因为字典序和日期顺序一致
  • 建议手动指定Schema,避免Spark自动推断Schema时出现错误

内容的提问来源于stack exchange,提问作者Md Sirajus Salayhin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:02:46