如何在AWS Glue中通过Workflow动态创建ETL作业及处理相关场景
AWS Glue动态/程序化ETL作业实现方案
1. S3新文件自动触发Glue作业
实现步骤
- 配置S3事件通知:
- 进入目标S3存储桶的「属性」->「事件通知」,创建新通知。
- 设置事件类型为
s3:ObjectCreated:*,可指定触发的文件前缀/后缀(例如仅触发.csv文件)。 - 选择目标为「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
- 进入AWS Glue控制台,选择「Workflows」->「Add workflow」,命名并创建。
- 添加触发器:选择「Add trigger」,类型为「Event bridge」,关联提前创建的CloudWatch Events规则(监听S3 ObjectCreated事件)。
- 添加作业:将动态ETL作业添加到Workflow中,配置作业参数引用触发器传递的S3文件信息。
- (可选)添加验证步骤:新增作业检查目标数据写入状态,确保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
相关产品推荐
相关产品推荐

