尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache Airflow 中 Amazon EMR 连接的配置与实战指南(EMR Connection)

Apache Airflow 中 Amazon EMR 连接的配置与实战指南(EMR Connection) Apache Airflow 中 Amazon EMR 连接的配置与实战指南EMR Connection【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的 Amazon provider 提供了一种专用的Amazon Elastic MapReduce (EMR) Connection连接类型用于集中存放启动 EMR 集群所需的初始参数并以 JSON 字典形式注入到 boto3 的run_job_flow调用中。本文以 providers/amazon/docs/connections/emr.rst 为核心骨架结合EmrHook源码、EmrCreateJobFlowOperator实现与单元测试系统讲解该连接的设计定位、配置方法、底层合并逻辑、常见报错排查以及在 DAG 中的落地用法帮助你正确区分EMR 集群配置连接与AWS 凭证连接避免参数校验失败。连接定位只存集群启动参数不存 AWS 凭证EMR Connection 在 Apache Airflow 的 Amazon provider 中承担着一个非常特殊的职责。官方文档开篇便明确指出This connection type is only used to store parameters to Start EMR Cluster (run_job_flowboto3 EMR client method).也就是说它不是用于保存 AWS 访问凭证的连接而是一个纯粹的参数仓库——所有写入该连接Extra字段的内容最终都会被透传给 boto3 EMR client 的run_job_flowAPI作为创建集群Job Flow的初始配置。这一点在源码中有清晰的体现。查看 EmrHook 实现class EmrHook(AwsBaseHook): conn_name_attr emr_conn_id default_conn_name emr_default conn_type emr hook_name Amazon Elastic MapReduce def __init__(self, emr_conn_id: str | None default_conn_name, *args, **kwargs) - None: self.emr_conn_id emr_conn_id kwargs[client_type] emr super().__init__(*args, **kwargs)EmrHook继承自AwsBaseHook其client_type被固定为emr而emr_conn_id是独立于aws_conn_id的连接标识。因此一个完整的 EMR 任务通常需要两个连接协同工作连接类型默认连接名职责Amazon Web Services Connectionaws_default提供 AWS 凭证Access Key / Session Token / IAM Role 等负责 AWS API 鉴权Amazon Elastic MapReduce Connectionemr_default仅提供run_job_flow的初始配置参数负责建什么样的集群AWS 凭证的配置方式请参见 Amazon Web Services Connection 文档本文不再展开。配置方法在 Extra 中填写 JSON 字典EMR Connection 的配置集中在Extra (optional)字段官方文档给出的要求如下Specify the parameters (as ajsondictionary) that can be used as an initial configuration inEmrHook.create_job_flowto propagate toRunJobFlow API. All parameters are optional.即将一组 JSON 字典形式的参数填入 Extra 字段作为EmrHook.create_job_flow()的初始配置最终原样传递给 AWSRunJobFlowAPI。所有参数均为可选参数名必须严格遵循 RunJobFlow API 的定义。Airflow 源码中为 EMR 连接注册了自定义的 UI 字段行为见 EmrHook.get_ui_field_behaviourclassmethod def get_ui_field_behaviour(cls) - dict[str, Any]: Return custom UI field behaviour for Amazon Elastic MapReduce Connection. return { hidden_fields: [host, schema, port, login, password], relabeling: { extra: Run Job Flow Configuration, }, placeholders: { extra: json.dumps( { Name: MyClusterName, ReleaseLabel: emr-5.36.0, Applications: [{Name: Spark}], Instances: { InstanceGroups: [ { Name: Primary node, Market: SPOT, InstanceRole: MASTER, InstanceType: m5.large, InstanceCount: 1, }, ], KeepJobFlowAliveWhenNoSteps: False, TerminationProtected: False, }, StepConcurrencyLevel: 2, }, indent2, ), }, }这段代码揭示了三个关键事实隐藏字段host、schema、port、login、password在连接编辑界面中全部隐藏因为该连接不需要任何主机/账号信息字段重命名界面上的Extra字段被重新标注为Run Job Flow Configuration明确提示其语义内置示例界面会给出一个 JSON 占位示例覆盖了集群名、EMR 发行版、应用栈、实例组、终止保护、步骤并发级别等最常用参数。一个可直接参考的 Extra 配置示例与 UI 占位符一致{ Name: MyClusterName, ReleaseLabel: emr-5.36.0, Applications: [ { Name: Spark } ], Instances: { InstanceGroups: [ { Name: Primary node, Market: SPOT, InstanceRole: MASTER, InstanceType: m5.large, InstanceCount: 1 } ], KeepJobFlowAliveWhenNoSteps: false, TerminationProtected: false }, StepConcurrencyLevel: 2 }常用字段速查均可选参数名与 AWS API 完全一致Name集群名称会显示在 EMR 控制台ReleaseLabelEMR 发行版标签例如emr-5.36.0、emr-6.15.0决定 Hadoop/Spark 等组件版本Applications要安装的应用列表如[{Name: Spark}]、[{Name: Hive}]Instances实例配置含InstanceGroups/InstanceFleets、KeepJobFlowAliveWhenNoSteps无步骤时是否保持集群存活、TerminationProtected终止保护等StepConcurrencyLevel集群可并行执行的步骤数量仅 EMR 5.28.0 及以上版本支持BootstrapActions引导操作用于在集群启动时安装自定义依赖Configurations应用级配置例如 Spark 内存参数、Hive 设置等SecurityConfiguration安全配置名称用于启用加密等能力JobFlowRole/ServiceRoleEC2 实例角色与 EMR 服务角色 ARNTags集群标签列表。参数校验失败不要在 Extra 里塞无关参数由于 Extra 内容最终会被原样展开为run_job_flow(**config)的关键字参数任何不属于RunJobFlowAPI 定义的参数都会触发 boto3 的参数校验错误。官方文档给出了典型报错Parameter validation failed: Unknown parameter in input: region_name, must be one of: ...例如把 AWS 区域region_name、凭证等看似合理的参数写进 EMR 连接的 Extra就会得到上述报错。原因是 boto3 客户端在调用run_job_flow前会做严格的参数白名单校验未知参数一律拒绝。正确的做法是区域region通过 AWS Connection 或aws_conn_id/ hook 的region_name参数指定凭证credentials通过 AWS Connection 指定EMR 连接 Extra只放RunJobFlowAPI 允许的参数。底层原理create_job_flow 的配置合并流程EmrHook.create_job_flow()是 EMR 连接与 AWS API 之间的桥梁其完整逻辑见 create_job_flow 实现def create_job_flow(self, job_flow_overrides: dict[str, Any]) - dict[str, Any]: config {} if self.emr_conn_id: try: emr_conn self.get_connection(self.emr_conn_id) except AirflowNotFoundException: warnings.warn( fUnable to find {self.hook_name} Connection ID {self.emr_conn_id!r}, using an empty initial configuration. ..., UserWarning, stacklevel2, ) else: if emr_conn.conn_type and emr_conn.conn_type ! self.conn_type: warnings.warn( f{self.hook_name} Connection expected connection type {self.conn_type!r}, ..., UserWarning, stacklevel2, ) config emr_conn.extra_dejson.copy() config.update(job_flow_overrides) response self.get_conn().run_job_flow(**config) return response从源码可以提炼出完整的执行链路读取初始配置通过emr_conn_id查找 EMR 连接将其Extra字段解析为 JSON 字典extra_dejson连接不存在时优雅降级若连接不存在抛出UserWarningUnable to find Amazon Elastic MapReduce Connection ID ...并以空字典作为初始配置继续执行连接类型校验若连接的conn_type不是emr会给出警告This connection might not work correctly.但不会中断执行覆盖合并config.update(job_flow_overrides)将 DAG 中传入的job_flow_overrides合并进初始配置同名键以 overrides 为准实现连接提供基线、任务按需覆盖发起调用最终以run_job_flow(**config)创建集群并返回 boto3 响应包含JobFlowId。上述第 2、3 步的行为均有对应的单元测试覆盖见 test_emr.pytest_empty_emr_conn_idemr_conn_idNone时run_job_flow仅收到job_flow_overridestest_missing_emr_conn_id连接不存在时发出UserWarning并继续执行test_emr_conn_id_wrong_conn_type连接类型错误时发出警告但调用仍成功。另外test_create_job_flow_extra_argstest_emr.py#L158-L183专门验证了EMR 新增 API 参数时无需改动 Airflow 代码的特性——只要参数合法Extra 中任意RunJobFlow参数都能直接透传这正是该连接被设计为透传参数容器的根本原因。在 DAG 中落地EmrCreateJobFlowOperator 与覆盖机制在 DAG 中最常用的入口是EmrCreateJobFlowOperator见 operators/emr.py 中的定义。它的核心参数包括emr_conn_id默认emr_default用于读取初始集群配置为None或连接不存在时使用空配置job_flow_overridesboto3 风格的参数字典或指向.json文件的路径支持模板渲染用于覆盖 EMR 连接中的同名参数aws_conn_idAWS 凭证连接为None时使用默认 boto3 凭证链region_nameAWS 区域未指定时使用默认 boto3 行为wait_for_completion/wait_policy控制任务是否等待集群就绪或步骤完成deferrable是否以 deferrable异步模式运行需要安装aiobotocoreterminate_job_flow_on_failure创建失败时是否尽力终止已创建的集群默认True。一个典型的 DAG 片段如下from airflow.providers.amazon.aws.operators.emr import EmrCreateJobFlowOperator create_emr_cluster EmrCreateJobFlowOperator( task_idcreate_emr_cluster, emr_conn_idemr_default, # 从 EMR 连接读取初始配置 aws_conn_idaws_default, # AWS 凭证连接 job_flow_overrides{ Name: MyAnalysisCluster, # 覆盖连接中的集群名 Instances: { InstanceGroups: [ { Name: Core node, Market: ON_DEMAND, InstanceRole: CORE, InstanceType: m5.xlarge, InstanceCount: 2, } ], KeepJobFlowAliveWhenNoSteps: True, }, }, )从 execute 方法 可以看到job_flow_overrides支持字符串形式通过ast.literal_eval解析调用hook.create_job_flow(job_flow_overrides)后校验 HTTP 状态码并记录JobFlowId随后通过EmrClusterLink/EmrLogsLink为任务生成指向 EMR 控制台的链接。在创建集群之后通常还会配合同一文件中的EmrAddStepsOperator向集群追加步骤、EmrTerminateJobFlowOperator终止集群等算子完成创建 → 提交作业 → 释放资源的完整生命周期。当EmrAddStepsOperator传入wait_for_completionTrue时底层会调用 EmrHook.add_job_flow_steps 并通过step_completewaiter 轮询步骤状态失败时还会进一步调用describe_step输出FailureDetails便于排查。为什么该连接无法测试连接在 Airflow UI 中对 EMR 连接点击 Test 会直接返回失败状态。这是有意为之见 test_connection 实现def test_connection(self): msg ( f{self.hook_name!r} Airflow Connection cannot be tested, by design it stores fonly key/value pairs and does not make a connection to an external resource. ) return False, msg源码注释说明了原因该 hook 基于AwsGenericHook若不做覆盖会默认使用 boto3 凭证策略去探测 AWS STS而这并非本连接的设计用途。EMR 连接按设计只存储键值对、不建立任何外部连接因此连接测试恒为失败。要验证 AWS 连通性请对 AWS Connection 执行连接测试。配置与管理建议结合源码行为与测试证据总结以下实践建议职责分离EMR 连接只放RunJobFlow参数集群名、发行版、实例组、应用栈等凭证与区域走 AWS Connection。这样同一套集群基线可被多个 DAG 复用。善用覆盖机制把稳定不变的基线如实例类型、发行版沉淀在连接 Extra 中把每次运行不同的内容如集群名放在job_flow_overrides里通过config.update(job_flow_overrides)的合并顺序实现灵活覆盖。参数名必须与 AWS API 对齐Extra 中的每个键都会被展开为run_job_flow的入参出现Unknown parameter in input报错时优先检查是否有region_name、access_key之类的非 API 参数混入。连接命名规范保持默认连接名emr_default可减少显式传参若自定义连接名需在每个使用EmrHook/EmrCreateJobFlowOperator的任务中显式指定emr_conn_id。不要依赖连接测试EMR 连接无法测试属于正常行为AWS 侧的连通性验证应通过 AWS Connection 完成。敏感信息不入 Extrarun_job_flow参数中不应包含明文密钥需要加密能力时应使用SecurityConfiguration或交由 AWS 服务端处理。通过正确配置 EMR Connection你可以将集群创建参数与 AWS 凭证彻底解耦让 DAG 更加简洁、可维护同时最大化复用既有的 AWS 基础设施与运维约定。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表