如何通过Glue Job逐个加载S3中的CSV文件至Redshift
实现Glue Job单次读取单个S3 CSV文件并加载到Redshift
方案一:参数化Glue Job,手动指定单个文件路径
适合按需指定特定文件处理的场景,每次运行Job时传入目标文件的S3路径:
配置Glue Job参数
在Glue Job的「Job parameters」中添加自定义参数--s3_file_path,值设为单个CSV文件的完整S3路径(例如s3://your-bucket/path/to/file1.csv)。在Glue脚本中获取参数并处理
用Glue工具类获取参数后直接读取单个文件,再写入Redshift:import sys from awsglue.utils import getResolvedOptions from pyspark.sql import SparkSession # 获取Job传入的参数 args = getResolvedOptions(sys.argv, ['s3_file_path', 'JOB_NAME']) s3_file_path = args['s3_file_path'] spark = SparkSession.builder.appName(args['JOB_NAME']).getOrCreate() # 读取单个CSV文件(根据实际格式调整header、分隔符等选项) df = spark.read.csv( s3_file_path, header=True, inferSchema=True, delimiter=',' ) # 加载数据到Redshift df.write \ .format("jdbc") \ .option("url", "jdbc:redshift://your-redshift-cluster:5439/your-db") \ .option("dbtable", "your-schema.your-target-table") \ .option("user", "your-redshift-username") \ .option("password", "your-redshift-password") \ .mode("append") # 按需选择append/overwrite模式 .save()运行Job
每次执行时修改--s3_file_path的值,即可切换处理不同的单个文件。
方案二:自动遍历S3文件,逐个触发Glue Job
如果需要批量处理所有文件但每次仅运行一个Job实例,可以用Lambda或Step Functions实现自动化:
获取S3目标目录下的CSV文件列表
用Boto3列出指定前缀下的所有CSV文件:import boto3 s3_client = boto3.client('s3') bucket = 'your-bucket-name' prefix = 'path/to/csv-files/' response = s3_client.list_objects_v2(Bucket=bucket, Prefix=prefix) csv_files = [obj['Key'] for obj in response['Contents'] if obj['Key'].endswith('.csv')]逐个触发Glue Job
遍历文件列表,每次调用Glue API触发Job并传入单个文件路径:glue_client = boto3.client('glue') for file_key in csv_files: full_s3_path = f's3://{bucket}/{file_key}' glue_client.start_job_run( JobName='your-glue-job-name', Arguments={ '--s3_file_path': full_s3_path } ) # 可选:添加延迟避免并发过高 import time time.sleep(30)
方案三:脚本内指定文件名过滤
如果不需要参数化,也可以直接在脚本中通过文件名匹配读取单个文件:
# 仅读取指定文件名的CSV文件 df = spark.read.csv( "s3://your-bucket/path/to/files/", header=True, inferSchema=True, pathGlobFilter="specific_file.csv" # 精确指定单个文件名 )
注意事项
- Schema一致性:确保所有CSV文件的列结构完全一致,避免写入Redshift时出现字段不匹配错误。
- 权限配置:Glue Job角色需要具备S3读取权限、Redshift写入权限,以及相关加密资源的访问权限(如果启用加密)。
- 写入模式:根据业务需求选择
append(追加数据)或overwrite(覆盖表)模式,避免数据重复或丢失。
内容的提问来源于stack exchange,提问作者Srikanth
相关产品推荐
相关产品推荐

