AWS Glue动态路径S3 CSV同步至Redshift:Catalog配置与脚本实现
解决AWS Glue同步S3动态日期路径到Redshift的问题
我刚好处理过类似的动态分区同步场景,下面分两部分帮你搞定这个需求:
一、在Glue数据目录中添加动态S3路径(仅最新日期文件夹)
你的S3路径是按d=YYYY-MM-DD分区的,我们可以利用Glue的分区表特性来管理,同时确保只同步最新的分区:
方法1:手动创建分区表 + 动态注册最新分区
先创建基础Glue表
- 登录Glue控制台,进入数据目录 → 表 → 添加表
- 数据源选S3,设置表的根路径为
s3://data-dl/abc/ - 定义分区键:添加名为
d的分区列,类型选字符串(格式为YYYY-MM-DD) - 配置CSV格式:设置分隔符、是否包含表头(如果你的CSV有表头,记得勾选“跳过表头行”),完成表创建。
动态注册最新分区到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
相关产品推荐
相关产品推荐

