Apache Airflow 集成 Amazon EMR:Job Flow 生命周期编排操作符与传感器完整指南
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Amazon EMR(前身为 Amazon Elastic MapReduce)是 AWS 提供的托管集群平台,用于在云上运行 Apache Hadoop、Apache Spark 等大数据框架,以处理和分析海量数据,并支持在 Amazon S3、Amazon DynamoDB 等 AWS 数据存储与数据库之间转换和搬运数据。本文以 Apache Airflow 官方仓库中的 emr.rst 为骨架,结合apache-airflow-providers-amazon的实际源码与系统测试示例,系统讲解如何使用 EMR 系列操作符(Operator)与传感器(Sensor)对 EMR Job Flow 实施全生命周期管理:从创建集群、添加步骤、修改集群、执行 Notebook,到终止集群与状态等待,并深入说明wait_policy、可延迟(deferrable)执行、OpenLineage 血缘注入等进阶能力。
读完本文,你将能够基于仓库提供的系统测试示例,独立编写一套可运行的 EMR 编排 DAG,并理解各操作符底层如何调用 Boto3 EMR API、如何配置等待策略与重试参数以规避 AWS 服务配额导致的限流。
前置准备:运行 EMR 操作符的前提条件
根据仓库中的 prerequisite_tasks.rst,使用 EMR 相关操作符前需要完成三件事:
通过 AWS Console 或 AWS CLI 创建所需 AWS 资源,其中最关键的是 EMR 所需的 IAM 服务角色。文档特别强调,要成功运行示例,必须创建
EMR_EC2_DefaultRole与EMR_DefaultRole两个 IAM 服务角色,一条命令即可完成:aws emr create-default-roles通过 pip 安装 Amazon 提供方包:
pip install 'apache-airflow[amazon]'配置 AWS 连接(Connection),详见仓库 connections/aws 文档。
通用参数:所有 EMR 操作符共享的 AWS 连接参数
EMR 系列操作符均继承自AwsBaseOperator,因此共享一组通用参数。其完整说明位于 generic_parameters.rst,要点如下:
| 参数 | 默认值 | 说明 |
|---|---|---|
aws_conn_id | aws_default | 引用的 AWS 连接 ID。若设为None,则跳过连接查找,直接使用默认 Boto3 行为(如环境变量凭证) |
region_name | None | AWS 区域名。为None或省略时,使用 AWS 连接 Extra 参数中的region_name |
verify | None | 是否校验 SSL 证书:False表示不校验;也可传 CA 证书 bundle 的文件路径。为None时使用连接 Extra 参数中的配置 |
botocore_config | None | 用于构造botocore.config.Config的字典,可配置重试策略、超时等,例如: |
{ "signature_version": "unsigned", "s3": { "us_east_1_regional_endpoint": True, }, "retries": { "mode": "standard", "max_attempts": 10, }, "connect_timeout": 300, "read_timeout": 300, "tcp_keepalive": True, }botocore_config是缓解 EMR 限流问题的关键手段(详见下文“限流与重试”一节)。注意:传入空字典{}会覆盖(而非合并)连接级配置。
创建 EMR Job Flow:EmrCreateJobFlowOperator
核心用法与 wait_policy 等待策略
使用 EmrCreateJobFlowOperator 可以创建新的 EMR Job Flow。集群通常在步骤执行完毕后由 EMR 自动终止。
该操作符默认行为是:集群一启动即把 DAG 任务节点标记为成功(wait_policy=None)。通过设置不同的wait_policy可以改变这一行为:
WaitPolicy.WAIT_FOR_COMPLETION—— DAG 任务节点等待集群进入Running状态;WaitPolicy.WAIT_FOR_STEPS_COMPLETION—— DAG 任务节点等待集群终止(即所有步骤完成后集群被关闭)。
从 utils/waiter.py 源码可见,WaitPolicy是一个字符串枚举,并通过WAITER_POLICY_NAME_MAPPING映射到对应的 Boto3 waiter:
class WaitPolicy(str, Enum): WAIT_FOR_COMPLETION = "wait_for_completion" # 映射到 waiter: job_flow_waiting WAIT_FOR_STEPS_COMPLETION = "wait_for_steps_completion" # 映射到 waiter: job_flow_terminated在 EmrCreateJobFlowOperator 的初始化逻辑中可以看到完整的参数处理:wait_policy与旧的wait_for_completion参数并存,二者同时出现时会给出弃用警告;不指定任何等待参数时,任务在集群创建成功后立即返回。
该操作符还提供了两个与等待行为相关的参数:
waiter_delay:轮询间隔(秒),默认60;waiter_max_attempts:最大轮询次数,默认60;terminate_job_flow_on_failure:默认True,当集群创建后发生失败时,会尽力终止该 Job Flow(清理失败不会掩盖原始异常)。
deferrable 可延迟模式
EmrCreateJobFlowOperator支持通过deferrable=True参数进入可延迟执行模式。在 deferrable 模式下,任务会释放 worker 槽位,显著提升 Airflow 集群资源利用效率——但前提是你的部署中必须运行Airflow triggerer组件。从源码 emr.py 第 866 行可以看到,deferrable 模式下会挂起EmrCreateJobFlowTrigger,由 triggerer 异步轮询集群状态,完成后通过execute_complete回调恢复任务。
JobFlow 配置:完整示例
创建 Job Flow 需要提供 EMR 集群的完整配置。仓库系统测试 example_emr.py 中给出了标准的配置写法,这里我们创建一个名为PiCalc的单节点 EMR 集群,它仅包含一个用 Spark 计算 π 值的步骤calculate_pi:
SPARK_STEPS = [ { "Name": "calculate_pi", "ActionOnFailure": "CONTINUE", "HadoopJarStep": { "Jar": "command-runner.jar", "Args": ["/usr/lib/spark/bin/run-example", "SparkPi", "10"], }, } ] JOB_FLOW_OVERRIDES: dict[str, Any] = { "Name": "PiCalc", "ReleaseLabel": "emr-7.1.0", "Applications": [{"Name": "Spark"}], "Instances": { "InstanceGroups": [ { "Name": "Primary node", "Market": "ON_DEMAND", "InstanceRole": "MASTER", "InstanceType": "m5.xlarge", "InstanceCount": 1, }, ], "KeepJobFlowAliveWhenNoSteps": True, "TerminationProtected": False, }, "Steps": SPARK_STEPS, "JobFlowRole": "EMR_EC2_DefaultRole", "ServiceRole": "EMR_DefaultRole", }配置要点说明:
'KeepJobFlowAliveWhenNoSteps': False会告诉集群在步骤执行完毕后自动关闭;若设置为True,则集群在步骤完成后保持存活,方便后续通过EmrAddStepsOperator继续追加步骤(系统测试正是利用这一点在同一个 DAG 里先创建、再修改、再追加步骤)。- 配置中也可以省略
Steps,之后再用EmrAddStepsOperator在任意时间点添加步骤(详见下一节)。 - 集群配置字段与 Boto3 EMR 客户端的
run_job_flow请求体一一对应,更多字段说明可查阅 Boto3 EMR 客户端文档。
提示:通过 EMR API(如本例)启动的集群默认对用户不可见,你可能会在 EMR 管理控制台中看不到它。若希望集群对所有用户可见,可在
JOB_FLOW_OVERRIDES字典末尾追加'VisibleToAllUsers': True。
创建 Job Flow 的 DAG 写法
create_job_flow = EmrCreateJobFlowOperator( task_id="create_job_flow", job_flow_overrides=JOB_FLOW_OVERRIDES, )该操作符会把新集群的JobFlowId作为返回值(output)传递给下游任务,供后续添加步骤、终止集群等操作直接使用——这正是系统测试 DAG 中modify_cluster、add_steps、remove_cluster均以create_job_flow.output作为job_flow_id的原因。此外,它还会通过EmrClusterLink与EmrLogsLink为任务注入集群控制台与日志的快捷链接。
添加步骤到已有集群:EmrAddStepsOperator
对于已存在的 Job Flow,使用 EmrAddStepsOperator 向其中添加步骤,它同样支持deferrable=True可延迟模式(deferrable 模式下会自动等待所有步骤完成,且wait_for_completion参数不生效)。
add_steps = EmrAddStepsOperator( task_id="add_steps", job_flow_id=create_job_flow.output, steps=SPARK_STEPS, execution_role_arn=execution_role_arn, ) add_steps.wait_for_completion = True参数细节:
job_flow_id与job_flow_name必须二选一:源码 emr.py 第 143 行通过exactly_one校验强制这一约束;使用job_flow_name时,操作符会调用hook.get_cluster_id_by_name在指定状态(cluster_states)的集群中按名称查找唯一匹配的集群 ID。steps:Boto3 风格的步骤列表,也支持传入.json文件路径(template_ext含.json),文件内容会被解析为步骤列表;steps字段本身可被模板渲染。wait_for_completion:默认False;为True时等待所有步骤完成(配合waiter_delay/waiter_max_attempts控制轮询频率与次数,默认分别为30与60)。execution_role_arn:步骤在集群上使用的运行时角色 ARN(即 EMR Runtime Role 场景)。- 返回值是新增步骤的 ID 列表,可将其中的步骤 ID 交给
EmrStepSensor监控,例如系统测试中的get_step_id(add_steps.output)。
OpenLineage 父作业信息注入
对于通过command-runner.jar以spark-submit或run-example方式启动的 Spark 步骤,EmrAddStepsOperator可以自动注入 OpenLineage 父作业信息,从而把 Spark 应用发出的 OpenLineage 事件关联回提交该步骤的 Airflow 任务。
启用方式有两种:
- 全局开启:设置 Airflow 配置项
[openlineage] spark_inject_parent_job_info; - 单操作符开启:传入
openlineage_inject_parent_job_info=True。
从源码 emr.py 第 159 行的_inject_openlineage_parent_job_information可以看到注入逻辑的细节:它只处理 Jar 名为command-runner.jar、且首参为run-example或spark-submit的步骤;若步骤参数中已存在spark.openlineage.parent*属性则跳过,避免覆盖用户手工配置;非 Spark 步骤完全不受影响。当然,注入的前提是 Spark 应用本身已安装并启用了 OpenLineage Spark 集成。
终止 EMR Job Flow:EmrTerminateJobFlowOperator
使用 EmrTerminateJobFlowOperator 终止一个 Job Flow,同样支持deferrable=True。
remove_cluster = EmrTerminateJobFlowOperator( task_id="remove_cluster", job_flow_id=create_job_flow.output, )系统测试中在终止任务上设置了trigger_rule = TriggerRule.ALL_DONE,确保无论前面的步骤是否失败,集群都会被清理,避免资源泄漏。EmrTerminateJobFlowOperator底层调用 EMRterminate_job_flowsAPI;deferrable 模式下由EmrTerminateJobFlowTrigger异步等待集群进入TERMINATED状态。
修改 EMR 集群:EmrModifyClusterOperator
文档示例(example_emr.py 第 155 行)演示了如何通过 EmrModifyClusterOperator 修改集群配置——注意该小节标题为 “Modify Amazon EMR container”,但示例代码实际使用的是修改 EMR 集群的EmrModifyClusterOperator,而非传感器:
modify_cluster = EmrModifyClusterOperator( task_id="modify_cluster", cluster_id=create_job_flow.output, step_concurrency_level=1 )该操作符用于调整运行中集群的运行时参数,示例中把step_concurrency_level(步骤并发级别)设置为1,使集群同一时间只执行一个步骤。这对于控制多个并发步骤的资源占用非常实用。
启动与停止 EMR Notebook 执行
对于挂载在运行中集群上的 EMR Notebook,可以编程方式启动/停止其执行。
启动执行:EmrStartNotebookExecutionOperator
EmrStartNotebookExecutionOperator 在指定 Notebook Editor 上启动一次执行,系统测试 example_emr_notebook_execution.py 中的用法:
start_execution = EmrStartNotebookExecutionOperator( task_id="start_execution", editor_id=editor_id, cluster_id=cluster_id, relative_path="EMR-System-Test.ipynb", service_role="EMR_Notebooks_DefaultRole", )editor_id:目标 Notebook Editor(EMR Notebook)的 ID;cluster_id:执行所使用的运行中集群 ID;relative_path:Notebook 文件在 Editor 中的相对路径(如EMR-System-Test.ipynb);service_role:执行 Notebook 所用的服务角色,示例使用EMR_Notebooks_DefaultRole。
执行启动后,操作符返回NotebookExecutionId,可传递给传感器监控其状态。
停止执行:EmrStopNotebookExecutionOperator
EmrStopNotebookExecutionOperator 停止一次正在运行的 Notebook 执行:
stop_execution = EmrStopNotebookExecutionOperator( task_id="stop_execution", notebook_execution_id=notebook_execution_id_1, )传感器(Sensors):状态等待与失败检测
EMR 传感器均继承自EmrBaseSensor,通过轮询 AWS API 获取状态,在达到目标状态时成功、进入失败状态时抛出异常,并支持 deferrable 模式。
EmrNotebookExecutionSensor:等待 Notebook 执行状态
EmrNotebookExecutionSensor 轮询describe_notebook_execution,默认目标状态为FINISHED(COMPLETED_STATES),失败状态为FAILED(FAILURE_STATES),也可通过target_states/failed_states自定义:
wait_for_execution_start = EmrNotebookExecutionSensor( task_id="wait_for_execution_start", notebook_execution_id=notebook_execution_id_1, target_states={"RUNNING"}, poke_interval=5, )系统测试正是用target_states={"RUNNING"}等待执行真正开始,再用target_states={"STOPPED"}确认停止动作完成,最后用默认状态FINISHED等待执行正常结束——展示了同一传感器在不同阶段的复用方式。
EmrJobFlowSensor:等待集群状态
EmrJobFlowSensor 轮询describe_cluster检查集群状态。源码中的默认值非常关键:
target_states默认["TERMINATED"]——默认等待集群被终止;failed_states默认["TERMINATED_WITH_ERRORS"];max_attempts默认60。
文档特别指出:将target_states设为['RUNNING', 'WAITING']可以等待集群真正就绪(即越过STARTING与BOOTSTRAPPING阶段):
check_job_flow = EmrJobFlowSensor(task_id="check_job_flow", job_flow_id=create_job_flow.output) check_job_flow.poke_interval = 10传感器会从Cluster.Status.State提取状态,并在失败时从StateChangeReason中提取失败码与错误信息用于诊断;同时它还会为任务挂载集群与日志快捷链接。deferrable 模式下,它挂起EmrTerminateJobFlowTrigger由 triggerer 异步等待。
EmrStepSensor:等待步骤状态
EmrStepSensor 监控单个步骤的执行状态,需要同时提供job_flow_id与step_id,默认等待步骤完成:
wait_for_step = EmrStepSensor( task_id="wait_for_step", job_flow_id=create_job_flow.output, step_id=get_step_id(add_steps.output), )限流与重试:应对 EMR 较低的服务配额
Amazon EMR 的服务配额相对较低,使用本页列出的任一操作符或传感器时都可能遇到throttling(限流)问题。文档给出的缓解思路是:自定义 AWS 连接配置,修改默认 Boto3 重试策略。
具体做法是在 AWS 连接的 Extra 参数(或操作符的botocore_config参数)中配置 botocore 的重试模式与次数,例如前文通用参数表中的配置:
{ "retries": { "mode": "standard", "max_attempts": 10, }, }将max_attempts从默认值调高(如10),配合standard模式的重试机制,可显著降低因瞬时限流导致的任务失败概率。若需要更细粒度的控制,还可以结合connect_timeout、read_timeout、tcp_keepalive等连接参数一并调整。
一个完整的端到端 EMR 编排示例
综合以上内容,可将仓库系统测试 example_emr.py 的 DAG 主体逻辑整理为如下任务链:创建安全配置 → 创建集群 → 修改集群并发级别 → 追加步骤并等待完成 → 监控步骤状态 → 终止集群 → 确认集群已终止:
from airflow.providers.amazon.aws.operators.emr import ( EmrAddStepsOperator, EmrCreateJobFlowOperator, EmrModifyClusterOperator, EmrTerminateJobFlowOperator, ) from airflow.providers.amazon.aws.sensors.emr import EmrJobFlowSensor, EmrStepSensor create_job_flow = EmrCreateJobFlowOperator( task_id="create_job_flow", job_flow_overrides=JOB_FLOW_OVERRIDES, ) modify_cluster = EmrModifyClusterOperator( task_id="modify_cluster", cluster_id=create_job_flow.output, step_concurrency_level=1 ) add_steps = EmrAddStepsOperator( task_id="add_steps", job_flow_id=create_job_flow.output, steps=SPARK_STEPS, ) add_steps.wait_for_completion = True wait_for_step = EmrStepSensor( task_id="wait_for_step", job_flow_id=create_job_flow.output, step_id=get_step_id(add_steps.output), ) remove_cluster = EmrTerminateJobFlowOperator( task_id="remove_cluster", job_flow_id=create_job_flow.output, ) remove_cluster.trigger_rule = TriggerRule.ALL_DONE check_job_flow = EmrJobFlowSensor(task_id="check_job_flow", job_flow_id=create_job_flow.output) check_job_flow.poke_interval = 10其中JOB_FLOW_OVERRIDES中的KeepJobFlowAliveWhenNoSteps设为True是关键:它保证集群在首个步骤完成后不被立即销毁,为后续“修改集群、追加步骤、步骤监控”等操作留出执行窗口,最后再由终止任务显式关闭集群——这正是 EMR 生命周期编排的典型模式。
参考资源
- 本页文档原文:providers/amazon/docs/operators/emr/emr.rst
- 操作符源码:providers/amazon/src/airflow/providers/amazon/aws/operators/emr.py(含
EmrCreateJobFlowOperator、EmrAddStepsOperator、EmrTerminateJobFlowOperator、EmrModifyClusterOperator、EmrStartNotebookExecutionOperator、EmrStopNotebookExecutionOperator) - 传感器源码:providers/amazon/src/airflow/providers/amazon/aws/sensors/emr.py(含
EmrNotebookExecutionSensor、EmrJobFlowSensor、EmrStepSensor) - 等待策略与 Waiter 映射:providers/amazon/src/airflow/providers/amazon/aws/utils/waiter.py
- 可运行的系统测试示例:
- example_emr.py(集群创建/修改/加步骤/终止/监控)
- example_emr_notebook_execution.py(Notebook 启动/停止/监控)
- 通用参数说明:providers/amazon/docs/_partials/generic_parameters.rst
- EMR 相关 Boto3 API 细节(
run_job_flow、add_job_flow_steps、describe_cluster、describe_notebook_execution等)可参考 AWS Boto3 EMR 客户端文档;IAM 角色创建命令参见 AWS CLI 的emr create-default-roles,角色配置详见 AWS EMR 管理指南中的 IAM 服务角色章节。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考