如何将Lambda事件获取的S3对象名传递给AWS Glue工作流及作业?
解决方案
方案1:批量传递对象列表,在Glue Job中循环处理所有文件
这种方式适合文件数量不多、单文件体积不大的场景,一次性把所有S3对象名传给Workflow,再由Job批量处理。
步骤1:Lambda中传递对象名到Glue Workflow
在Lambda代码里,先解析S3 Put事件拿到所有对象名(注意:S3事件可能一次包含多个Put记录),然后通过boto3启动Glue Workflow时,将对象名列表以JSON字符串的形式存入Workflow Run Properties:
import boto3 import json glue_client = boto3.client('glue') def lambda_handler(event, context): # 从S3事件中提取所有上传的XLSX对象名 object_names = [record['s3']['object']['key'] for record in event['Records'] if record['s3']['object']['key'].endswith('.xlsx')] if not object_names: return {'statusCode': 200, 'body': 'No XLSX files to process'} # 启动Glue Workflow,传递对象名列表作为运行属性 response = glue_client.start_workflow_run( Name='YourWorkflowName', RunProperties={ 'OBJECT_NAMES': json.dumps(object_names) } ) return {'statusCode': 200, 'body': f"Workflow started with run ID: {response['RunId']}"}
步骤2:Glue Job中读取并处理所有对象
在Glue Job的Python脚本里,先从Workflow的运行属性中取出对象名列表,然后循环处理每个XLSX文件:
import sys import json from awsglue.utils import getResolvedOptions from pyspark.sql import SparkSession # 获取Workflow传递的参数 args = getResolvedOptions(sys.argv, ['JOB_NAME', 'OBJECT_NAMES']) object_names = json.loads(args['OBJECT_NAMES']) spark = SparkSession.builder.appName(args['JOB_NAME']).getOrCreate() # 循环处理每个XLSX文件 for obj_key in object_names: s3_input_path = f"s3://your-bucket-name/{obj_key}" # 生成CSV输出路径(替换后缀为csv) s3_output_path = s3_input_path.replace('.xlsx', '.csv') # 读取XLSX文件(需要确保Glue Job安装了openpyxl依赖) df = spark.read.format("com.crealytics.spark.excel") \ .option("header", "true") \ .load(s3_input_path) # 写入CSV文件 df.write.mode("overwrite").csv(s3_output_path, header=True) spark.stop()
注意:Glue Job需要安装Spark Excel依赖,可在Job的"Job parameters"里添加
--extra-jars s3://your-bucket/path/to/spark-excel.jar,或者使用Glue自定义镜像。
方案2:逐个传递对象名,Workflow中触发多次Job实例
如果需要对每个文件独立处理(比如大文件拆分、失败重试隔离),可以选择以下两种方式实现:
方式A:Lambda中逐个启动Workflow运行
直接在Lambda里循环每个对象名,分别启动Workflow,每次传递单个对象名:
import boto3 glue_client = boto3.client('glue') def lambda_handler(event, context): object_names = [record['s3']['object']['key'] for record in event['Records'] if record['s3']['object']['key'].endswith('.xlsx')] for obj_key in object_names: glue_client.start_workflow_run( Name='YourWorkflowName', RunProperties={ 'OBJECT_NAME': obj_key } ) return {'statusCode': 200, 'body': f"Started {len(object_names)} workflow runs"}
随后在Glue Job脚本中读取OBJECT_NAME参数,处理单个文件即可。
方式B:Workflow内部实现批量触发
若希望通过一次Workflow启动处理所有文件,可在Workflow中添加一个Glue Python Shell Job作为入口:该Job读取初始传递的对象列表,再通过boto3调用start_job_run逐个触发目标转换Job,并传递单个对象名。
关键注意点
- 权限配置:确保Lambda角色拥有
glue:StartWorkflowRun权限;Glue Job角色拥有s3:GetObject、s3:PutObject权限(若用方案2B,还需glue:StartJobRun权限)。 - 对象名解码:S3对象名含特殊字符时,需用
urllib.parse.unquote解码,避免路径错误。 - 失败处理:可为Glue Job配置重试机制,或在Workflow中添加失败告警触发器。
内容的提问来源于stack exchange,提问作者Srinivas
相关产品推荐
相关产品推荐

