Apache Airflow 集成 Amazon EMRJob Flow 生命周期编排操作符与传感器完整指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAmazon 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_idaws_default引用的 AWS 连接 ID。若设为None则跳过连接查找直接使用默认 Boto3 行为如环境变量凭证region_nameNoneAWS 区域名。为None或省略时使用 AWS 连接 Extra 参数中的region_nameverifyNone是否校验 SSL 证书False表示不校验也可传 CA 证书 bundle 的文件路径。为None时使用连接 Extra 参数中的配置botocore_configNone用于构造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 FlowEmrCreateJobFlowOperator核心用法与 wait_policy 等待策略使用 EmrCreateJobFlowOperator 可以创建新的 EMR Job Flow。集群通常在步骤执行完毕后由 EMR 自动终止。该操作符默认行为是集群一启动即把 DAG 任务节点标记为成功wait_policyNone。通过设置不同的wait_policy可以改变这一行为WaitPolicy.WAIT_FOR_COMPLETION—— DAG 任务节点等待集群进入Running状态WaitPolicy.WAIT_FOR_STEPS_COMPLETION—— DAG 任务节点等待集群终止即所有步骤完成后集群被关闭。从 utils/waiter.py 源码可见WaitPolicy是一个字符串枚举并通过WAITER_POLICY_NAME_MAPPING映射到对应的 Boto3 waiterclass 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轮询间隔秒默认60waiter_max_attempts最大轮询次数默认60terminate_job_flow_on_failure默认True当集群创建后发生失败时会尽力终止该 Job Flow清理失败不会掩盖原始异常。deferrable 可延迟模式EmrCreateJobFlowOperator支持通过deferrableTrue参数进入可延迟执行模式。在 deferrable 模式下任务会释放 worker 槽位显著提升 Airflow 集群资源利用效率——但前提是你的部署中必须运行Airflow triggerer组件。从源码 emr.py 第 866 行可以看到deferrable 模式下会挂起EmrCreateJobFlowTrigger由 triggerer 异步轮询集群状态完成后通过execute_complete回调恢复任务。JobFlow 配置完整示例创建 Job Flow 需要提供 EMR 集群的完整配置。仓库系统测试 example_emr.py 中给出了标准的配置写法这里我们创建一个名为PiCalc的单节点 EMR 集群它仅包含一个用 Spark 计算 π 值的步骤calculate_piSPARK_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_idcreate_job_flow, job_flow_overridesJOB_FLOW_OVERRIDES, )该操作符会把新集群的JobFlowId作为返回值output传递给下游任务供后续添加步骤、终止集群等操作直接使用——这正是系统测试 DAG 中modify_cluster、add_steps、remove_cluster均以create_job_flow.output作为job_flow_id的原因。此外它还会通过EmrClusterLink与EmrLogsLink为任务注入集群控制台与日志的快捷链接。添加步骤到已有集群EmrAddStepsOperator对于已存在的 Job Flow使用 EmrAddStepsOperator 向其中添加步骤它同样支持deferrableTrue可延迟模式deferrable 模式下会自动等待所有步骤完成且wait_for_completion参数不生效。add_steps EmrAddStepsOperator( task_idadd_steps, job_flow_idcreate_job_flow.output, stepsSPARK_STEPS, execution_role_arnexecution_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。stepsBoto3 风格的步骤列表也支持传入.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_infoTrue。从源码 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 FlowEmrTerminateJobFlowOperator使用 EmrTerminateJobFlowOperator 终止一个 Job Flow同样支持deferrableTrue。remove_cluster EmrTerminateJobFlowOperator( task_idremove_cluster, job_flow_idcreate_job_flow.output, )系统测试中在终止任务上设置了trigger_rule TriggerRule.ALL_DONE确保无论前面的步骤是否失败集群都会被清理避免资源泄漏。EmrTerminateJobFlowOperator底层调用 EMRterminate_job_flowsAPIdeferrable 模式下由EmrTerminateJobFlowTrigger异步等待集群进入TERMINATED状态。修改 EMR 集群EmrModifyClusterOperator文档示例example_emr.py 第 155 行演示了如何通过 EmrModifyClusterOperator 修改集群配置——注意该小节标题为 “Modify Amazon EMR container”但示例代码实际使用的是修改 EMR 集群的EmrModifyClusterOperator而非传感器modify_cluster EmrModifyClusterOperator( task_idmodify_cluster, cluster_idcreate_job_flow.output, step_concurrency_level1 )该操作符用于调整运行中集群的运行时参数示例中把step_concurrency_level步骤并发级别设置为1使集群同一时间只执行一个步骤。这对于控制多个并发步骤的资源占用非常实用。启动与停止 EMR Notebook 执行对于挂载在运行中集群上的 EMR Notebook可以编程方式启动/停止其执行。启动执行EmrStartNotebookExecutionOperatorEmrStartNotebookExecutionOperator 在指定 Notebook Editor 上启动一次执行系统测试 example_emr_notebook_execution.py 中的用法start_execution EmrStartNotebookExecutionOperator( task_idstart_execution, editor_ideditor_id, cluster_idcluster_id, relative_pathEMR-System-Test.ipynb, service_roleEMR_Notebooks_DefaultRole, )editor_id目标 Notebook EditorEMR Notebook的 IDcluster_id执行所使用的运行中集群 IDrelative_pathNotebook 文件在 Editor 中的相对路径如EMR-System-Test.ipynbservice_role执行 Notebook 所用的服务角色示例使用EMR_Notebooks_DefaultRole。执行启动后操作符返回NotebookExecutionId可传递给传感器监控其状态。停止执行EmrStopNotebookExecutionOperatorEmrStopNotebookExecutionOperator 停止一次正在运行的 Notebook 执行stop_execution EmrStopNotebookExecutionOperator( task_idstop_execution, notebook_execution_idnotebook_execution_id_1, )传感器Sensors状态等待与失败检测EMR 传感器均继承自EmrBaseSensor通过轮询 AWS API 获取状态在达到目标状态时成功、进入失败状态时抛出异常并支持 deferrable 模式。EmrNotebookExecutionSensor等待 Notebook 执行状态EmrNotebookExecutionSensor 轮询describe_notebook_execution默认目标状态为FINISHEDCOMPLETED_STATES失败状态为FAILEDFAILURE_STATES也可通过target_states/failed_states自定义wait_for_execution_start EmrNotebookExecutionSensor( task_idwait_for_execution_start, notebook_execution_idnotebook_execution_id_1, target_states{RUNNING}, poke_interval5, )系统测试正是用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_idcheck_job_flow, job_flow_idcreate_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_idwait_for_step, job_flow_idcreate_job_flow.output, step_idget_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_idcreate_job_flow, job_flow_overridesJOB_FLOW_OVERRIDES, ) modify_cluster EmrModifyClusterOperator( task_idmodify_cluster, cluster_idcreate_job_flow.output, step_concurrency_level1 ) add_steps EmrAddStepsOperator( task_idadd_steps, job_flow_idcreate_job_flow.output, stepsSPARK_STEPS, ) add_steps.wait_for_completion True wait_for_step EmrStepSensor( task_idwait_for_step, job_flow_idcreate_job_flow.output, step_idget_step_id(add_steps.output), ) remove_cluster EmrTerminateJobFlowOperator( task_idremove_cluster, job_flow_idcreate_job_flow.output, ) remove_cluster.trigger_rule TriggerRule.ALL_DONE check_job_flow EmrJobFlowSensor(task_idcheck_job_flow, job_flow_idcreate_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.pyNotebook 启动/停止/监控通用参数说明providers/amazon/docs/_partials/generic_parameters.rstEMR 相关 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),仅供参考