如何在基于原生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
相关产品推荐
相关产品推荐

