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

从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)两组任务并行执行:

  1. 步骤A执行完成后,才能启动B1和B2
  2. 步骤B1执行完成后,才能启动C1
  3. 步骤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任务:

  1. 定义状态机,包含以下状态:
    • 启动EMR集群状态
    • 执行步骤A的任务状态
    • 并行分支状态:分别执行(B1→C1)和(B2→C2)
    • 终止EMR集群状态
  2. 通过Lambda触发Step Functions状态机,由状态机负责调度整个流程

该方案适合需要精细化控制流程、集成其他AWS服务的场景,但实现复杂度高于方案1。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 11:03:10