从Lambda启动的临时AWS EMR并行化步骤的最佳实践
问题背景与优化需求
场景概述
当S3桶s3://data-bucket/raw/接收到新文件时,触发Lambda函数,由该函数启动临时AWS EMR集群,运行以下Spark处理流水线:
- 步骤A:预处理流水线
- 输入:
s3://data-bucket/raw/file.YYYY-MM-DD.HH-MM-SS.csv - 输出:
s3://data-bucket/preprocessed/file.YYYY-MM-DD.HH-MM-SS.csv
- 输入:
- 步骤B1:特征提取流水线1
- 输入:
s3://data-bucket/preprocessed/file.YYYY-MM-DD.HH-MM-SS.csv - 输出:
s3://data-bucket/features/pipeline1/file.YYYY-MM-DD.HH-MM-SS.csv
- 输入:
- 步骤C1:机器学习流水线1
- 输入:
s3://data-bucket/features/pipeline1/file.YYYY-MM-DD.HH-MM-SS.csv - 输出:
s3://data-bucket/predictions/pipeline1/file.YYYY-MM-DD.HH-MM-SS.csv
- 输入:
- 步骤B2:特征提取流水线2
- 输入:
s3://data-bucket/preprocessed/file.YYYY-MM-DD.HH-MM-SS.csv - 输出:
s3://data-bucket/features/pipeline2/file.YYYY-MM-DD.HH-MM-SS.csv
- 输入:
- 步骤C2:机器学习流水线2
- 输入:
s3://data-bucket/features/pipeline2/file.YYYY-MM-DD.HH-MM-SS.csv - 输出:
s3://data-bucket/predictions/pipeline2/file.YYYY-MM-DD.HH-MM-SS.csv
- 输入:
现有实现(顺序执行)
当前采用全顺序执行流程:A→B1→C1→B2→C2,核心代码如下:
import boto3 emr = boto3.client('emr') # 定义各步骤的spark-submit命令,示例结构 spark_submits = { "A": ["spark-submit", "--class", "com.example.Preprocess", "s3://code-bucket/preprocess.jar"], "B1": ["spark-submit", "--class", "com.example.FeatureExtraction1", "s3://code-bucket/feature1.jar"], "C1": ["spark-submit", "--class", "com.example.MLPipeline1", "s3://code-bucket/ml1.jar"], "B2": ["spark-submit", "--class", "com.example.FeatureExtraction2", "s3://code-bucket/feature2.jar"], "C2": ["spark-submit", "--class", "com.example.MLPipeline2", "s3://code-bucket/ml2.jar"] } steps = [ { 'Name': step_name, 'ActionOnFailure': 'TERMINATE_CLUSTER' if step_name == 'A' else 'CONTINUE', 'HadoopJarStep': { 'Jar': 's3://eu-west-1.elasticmapreduce/libs/script-runner/script-runner.jar', 'Args': [*spark_submit_cmd] } } for step_name, spark_submit_cmd in spark_submits.items() ] response = emr.run_job_flow( Name='MyJobFlow', Steps=steps, # 其他集群配置(实例类型、数量、IAM角色等) ReleaseLabel='emr-6.10.0', Instances={ 'InstanceGroups': [ { 'Name': 'Master', 'InstanceRole': 'MASTER', 'InstanceType': 'm5.xlarge', 'InstanceCount': 1 }, { 'Name': 'Core', 'InstanceRole': 'CORE', 'InstanceType': 'm5.xlarge', 'InstanceCount': 2 } ], 'KeepJobFlowAliveWhenNoSteps': False, 'TerminationProtected': False }, JobFlowRole='EMR_EC2_DefaultRole', ServiceRole='EMR_DefaultRole' )
该方式因所有步骤串行执行,整体耗时过长,需优化。
优化需求
需满足以下依赖关系,同时实现(B1→C1)与(B2→C2)两组任务并行执行:
- 步骤A执行完成后,才能启动B1和B2
- 步骤B1执行完成后,才能启动C1
- 步骤B2执行完成后,才能启动C2
禁止使用简陋的bash脚本实现依赖控制。
最佳实践方案
方案1:利用EMR原生步骤依赖(Step Dependencies)
EMR支持为每个步骤指定DependsOn参数,通过该参数可以定义步骤间的依赖关系,从而实现并行执行逻辑,无需额外工具。
修改后的代码示例:
import boto3 emr = boto3.client('emr') spark_submits = { "A": ["spark-submit", "--class", "com.example.Preprocess", "s3://code-bucket/preprocess.jar"], "B1": ["spark-submit", "--class", "com.example.FeatureExtraction1", "s3://code-bucket/feature1.jar"], "C1": ["spark-submit", "--class", "com.example.MLPipeline1", "s3://code-bucket/ml1.jar"], "B2": ["spark-submit", "--class", "com.example.FeatureExtraction2", "s3://code-bucket/feature2.jar"], "C2": ["spark-submit", "--class", "com.example.MLPipeline2", "s3://code-bucket/ml2.jar"] } # 定义带依赖的步骤列表 steps = [ # 步骤A:无依赖,失败则终止集群 { 'Name': 'Step-A-Preprocess', 'ActionOnFailure': 'TERMINATE_CLUSTER', 'HadoopJarStep': { 'Jar': 's3://eu-west-1.elasticmapreduce/libs/script-runner/script-runner.jar', 'Args': [*spark_submits["A"]] } }, # 步骤B1:依赖步骤A完成 { 'Name': 'Step-B1-FeatureExtraction1', 'ActionOnFailure': 'CONTINUE', 'DependsOn': ['Step-A-Preprocess'], 'HadoopJarStep': { 'Jar': 's3://eu-west-1.elasticmapreduce/libs/script-runner/script-runner.jar', 'Args': [*spark_submits["B1"]] } }, # 步骤C1:依赖步骤B1完成 { 'Name': 'Step-C1-MLPipeline1', 'ActionOnFailure': 'CONTINUE', 'DependsOn': ['Step-B1-FeatureExtraction1'], 'HadoopJarStep': { 'Jar': 's3://eu-west-1.elasticmapreduce/libs/script-runner/script-runner.jar', 'Args': [*spark_submits["C1"]] } }, # 步骤B2:依赖步骤A完成 { 'Name': 'Step-B2-FeatureExtraction2', 'ActionOnFailure': 'CONTINUE', 'DependsOn': ['Step-A-Preprocess'], 'HadoopJarStep': { 'Jar': 's3://eu-west-1.elasticmapreduce/libs/script-runner/script-runner.jar', 'Args': [*spark_submits["B2"]] } }, # 步骤C2:依赖步骤B2完成 { 'Name': 'Step-C2-MLPipeline2', 'ActionOnFailure': 'CONTINUE', 'DependsOn': ['Step-B2-FeatureExtraction2'], 'HadoopJarStep': { 'Jar': 's3://eu-west-1.elasticmapreduce/libs/script-runner/script-runner.jar', 'Args': [*spark_submits["C2"]] } } ] response = emr.run_job_flow( Name='MyParallelJobFlow', Steps=steps, # 其他集群配置保持不变 ReleaseLabel='emr-6.10.0', Instances={ 'InstanceGroups': [ { 'Name': 'Master', 'InstanceRole': 'MASTER', 'InstanceType': 'm5.xlarge', 'InstanceCount': 1 }, { 'Name': 'Core', 'InstanceRole': 'CORE', 'InstanceType': 'm5.xlarge', 'InstanceCount': 2 } ], 'KeepJobFlowAliveWhenNoSteps': False, 'TerminationProtected': False }, JobFlowRole='EMR_EC2_DefaultRole', ServiceRole='EMR_DefaultRole' )
关键说明:
DependsOn参数接收步骤名称的列表,指定当前步骤必须等待哪些步骤完成后才能启动- B1和B2均依赖A,因此A完成后两者会并行启动
- C1依赖B1、C2依赖B2,保证组内的串行关系
- 集群会自动处理步骤的调度,无需额外编写控制逻辑
方案2:AWS Step Functions 编排(适合复杂场景)
如果后续需要更复杂的流程控制(如重试、分支、错误处理),可以使用AWS Step Functions来编排EMR任务:
- 定义状态机,包含以下状态:
- 启动EMR集群状态
- 执行步骤A的任务状态
- 并行分支状态:分别执行(B1→C1)和(B2→C2)
- 终止EMR集群状态
- 通过Lambda触发Step Functions状态机,由状态机负责调度整个流程
该方案适合需要精细化控制流程、集成其他AWS服务的场景,但实现复杂度高于方案1。
内容的提问来源于stack exchange,提问作者Mistapopo
相关产品推荐
相关产品推荐

