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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 21:15:33