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

如何在AWS Glue中通过Workflow动态创建ETL作业及处理相关场景

AWS Glue动态/程序化ETL作业实现方案

1. S3新文件自动触发Glue作业

实现步骤

  • 配置S3事件通知:
    1. 进入目标S3存储桶的「属性」->「事件通知」,创建新通知。
    2. 设置事件类型为s3:ObjectCreated:*,可指定触发的文件前缀/后缀(例如仅触发.csv文件)。
    3. 选择目标为「Lambda函数」,关联提前创建好的Lambda函数。
  • 编写Lambda触发逻辑:
    用boto3调用Glue API启动作业,示例代码:
    import boto3
    import os
    
    glue_client = boto3.client('glue')
    GLUE_JOB_NAME = os.environ['GLUE_JOB_NAME']
    
    def lambda_handler(event, context):
        # 提取S3文件信息,传递给Glue作业作为参数
        s3_record = event['Records'][0]['s3']
        source_bucket = s3_record['bucket']['name']
        source_key = s3_record['object']['key']
    
        # 启动Glue作业
        response = glue_client.start_job_run(
            JobName=GLUE_JOB_NAME,
            Arguments={
                '--source_bucket': source_bucket,
                '--source_key': source_key,
                '--target_bucket': os.environ['TARGET_BUCKET']
            }
        )
        return {'job_run_id': response['JobRunId']}
    
  • 权限配置:
    给Lambda角色添加glue:StartJobRun权限,同时确保S3存储桶有权限触发该Lambda函数。

2. 借助配置文件实现动态映射与转换

配置文件设计(JSON格式,存储在S3)

{
  "source_format": "csv",
  "source_options": {"withHeader": true},
  "mapping_rules": [
    {"source_field": "user_id", "target_field": "user_identifier", "transform": "cast_int"},
    {"source_field": "signup_date", "target_field": "registration_date", "transform": "to_date('yyyy-MM-dd')"},
    {"source_field": "email", "target_field": "user_email", "transform": "lowercase"}
  ],
  "target_format": "parquet",
  "target_options": {"compression": "snappy"}
}

Glue作业中读取配置并动态处理

import boto3
import json
import sys
from awsglue.context import GlueContext
from awsglue.dynamicframe import DynamicFrame
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from pyspark.sql.functions import col, to_date, lower

sc = SparkContext.getOrCreate()
glueContext = GlueContext(sc)
spark = glueContext.spark_session

# 读取S3中的配置文件
s3 = boto3.client('s3')
config_bucket = 'your-config-bucket'
config_key = 'etl/config/mapping_rules.json'
config_obj = s3.get_object(Bucket=config_bucket, Key=config_key)
config = json.loads(config_obj['Body'].read().decode('utf-8'))

# 获取Lambda传递的作业参数
args = getResolvedOptions(sys.argv, ['source_bucket', 'source_key', 'target_bucket'])
source_path = f"s3://{args['source_bucket']}/{args['source_key']}"
target_path = f"s3://{args['target_bucket']}/processed/"

# 加载源数据
source_dyf = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": [source_path]},
    format=config['source_format'],
    format_options=config['source_options']
)

# 动态执行字段映射与转换
def apply_transforms(df, rules):
    transformed_df = df
    for rule in rules:
        src_col = col(rule['source_field'])
        # 根据配置执行转换逻辑
        if rule['transform'] == 'cast_int':
            transformed_df = transformed_df.withColumn(rule['target_field'], src_col.cast('int'))
        elif rule['transform'].startswith('to_date'):
            date_format = rule['transform'].split("'")[1]
            transformed_df = transformed_df.withColumn(rule['target_field'], to_date(src_col, date_format))
        elif rule['transform'] == 'lowercase':
            transformed_df = transformed_df.withColumn(rule['target_field'], lower(src_col))
    # 仅保留目标字段
    target_fields = [rule['target_field'] for rule in rules]
    return transformed_df.select(target_fields)

# 转换为Spark DataFrame处理,再转回DynamicFrame
source_df = source_dyf.toDF()
transformed_df = apply_transforms(source_df, config['mapping_rules'])
transformed_dyf = DynamicFrame.fromDF(transformed_df, glueContext, "transformed_dyf")

# 写入目标S3
glueContext.write_dynamic_frame.from_options(
    frame=transformed_dyf,
    connection_type="s3",
    connection_options={"path": target_path},
    format=config['target_format'],
    format_options=config['target_options']
)

3. 创建动态/程序化Glue Workflow

控制台手动创建Workflow

  1. 进入AWS Glue控制台,选择「Workflows」->「Add workflow」,命名并创建。
  2. 添加触发器:选择「Add trigger」,类型为「Event bridge」,关联提前创建的CloudWatch Events规则(监听S3 ObjectCreated事件)。
  3. 添加作业:将动态ETL作业添加到Workflow中,配置作业参数引用触发器传递的S3文件信息。
  4. (可选)添加验证步骤:新增作业检查目标数据写入状态,确保ETL链路完整性。

程序化创建Workflow(boto3示例)

import boto3

glue_client = boto3.client('glue')

# 创建Workflow
workflow_name = 'dynamic-etl-workflow'
glue_client.create_workflow(
    Name=workflow_name,
    Description='Dynamic ETL workflow triggered by S3 new files'
)

# 创建触发器(关联CloudWatch Events规则,需提前配置S3事件监听)
trigger_name = 's3-new-file-trigger'
glue_client.create_trigger(
    Name=trigger_name,
    Type='EVENT',
    WorkflowName=workflow_name,
    Actions=[{'JobName': 'your-dynamic-glue-job'}],
    EventBatchingCondition={'BatchSize': 1, 'BatchWindow': 5}
)

# 启动Workflow(测试用)
glue_client.start_workflow_run(Name=workflow_name)

核心注意点

  • Workflow中的作业通过参数传递实现动态性,可将S3路径、配置文件路径作为参数传入。
  • 利用Workflow的监控面板可追踪每个步骤的执行状态,快速定位故障节点。

内容的提问来源于stack exchange,提问作者gvrspk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 05:03:18