ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Airflow Amazon Provider 前置任务指南:从资源创建、安装到连接配置

Apache Airflow Amazon Provider 前置任务指南:从资源创建、安装到连接配置 Apache Airflow Amazon Provider 前置任务指南从资源创建、安装到连接配置【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读在 Apache Airflow 中使用 Amazon Providerapache-airflow-providers-amazon提供的各类 Operator、Hook 与 Sensor覆盖 S3、EC2、ECS、Athena、EMR、Redshift 等数十个 AWS 服务之前必须先完成一组统一的前置任务在 AWS 侧准备好所需资源、通过 pip 安装 Provider 包及其 API 依赖库、并在 Airflow 中配置可用的 AWS 连接。本篇指南以仓库中 providers/amazon/docs/_partials/prerequisite_tasks.rst 为骨架结合源码与官方文档完整讲解这三步的实操方法、底层原理与常见坑点读完后你可以为任意一个 Amazon Provider 组件快速完成环境就绪。为什么会有“前置任务”这一节Amazon Provider 的几乎所有组件文档都要求在正文之前先完成统一的环境准备工作。这个约定通过 RST 的.. include::指令实现仓库根目录 providers/amazon/docs/_partials/prerequisite_tasks.rst 是一段被复用的共享片段被 operators、transfer、secrets-backends 等目录下的大量文档引用。例如以下文档都在对应位置引入了同一份前置任务片段EC2 操作指南.. include:: ../_partials/prerequisite_tasks.rstAthena SQL 操作指南.. include:: ../../_partials/prerequisite_tasks.rstAppFlow 操作指南Batch 操作指南这种“一次定义、多处复用”的文档组织方式保证了所有 AWS 组件的使用说明在环境准备环节口径完全一致也意味着只要完成本文介绍的三步就能平滑迁移到任意一个 Amazon Provider 组件的使用场景中。前置任务共三步使用AWS Console或AWS CLI创建必要资源通过pip安装 API 依赖库即安装 Provider 包在 Airflow 中Setup Connection配置 AWS 连接。下面逐一展开。第一步在 AWS 侧创建必要资源Amazon Provider 中的每个 Operator 都对应一个真实的 AWS 服务调用因此在本地运行任何 DAG 之前需要先在 AWS 账户中准备好目标资源。例如使用 S3 相关 Hook/Operator 前需要先创建 S3 Bucket使用 EC2 操作符前需要先准备好 EC2 实例或 AMI使用 EMR 操作符前需要先规划好 EMR 集群与 IAM 角色使用 Athena 前需要先在 S3 中准备数据并定义表结构。资源创建有两种途径AWS Console通过 Web 控制台交互式创建资源适合首次使用、可视化排查的场景AWS CLI通过命令行脚本化创建资源适合可重复执行的自动化场景。前提是本机已安装并配置好 AWS CLI 凭证aws configure。注意具体资源清单取决于你实际要使用的组件。仓库中每个组件的独立文档如 Athena 指南、EC2 指南都会先介绍该服务是什么再给出该服务的 Operator 用法示例资源准备应参照对应服务文档进行。第二步通过 pip 安装 Provider 包安装命令如下pip install apache-airflow[amazon]该命令会安装 Airflow 本体并同时拉取 Amazon Provider 包apache-airflow-providers-amazon及其全部核心 API 依赖库。关于安装包的依赖构成可以从仓库中的 providers/amazon/pyproject.toml 得到精确依据。以当前仓库版本provider 版本 9.36.0为例其核心依赖包括依赖说明apache-airflow2.11.0Provider 所要求的最低 Airflow 版本boto31.41.0、botocore1.41.0AWS SDK for Python是所有 AWS 组件与 AWS API 通信的基础redshift_connector2.1.3Redshift 连接驱动PyAthena3.10.0Athena 查询客户端watchtower3.3.1,4CloudWatch 日志集成jsonpath_ng、jmespath、inflection、marshmallow数据解析与序列化辅助库此外部分高级场景需要额外的可选依赖[project.optional-dependencies]段见 providers/amazon/pyproject.tomlaiobotocore使用异步 boto 会话deferrable 场景时安装s3fs使用 S3 文件系统如S3FileSystem时安装python3-saml使用 SAML 联邦认证时安装cncf.kubernetes、apache.hive、google等跨 Provider 协作场景安装。如需自定义安装子集也可以直接安装 Provider 本身例如pip install apache-airflow-providers-amazon[s3fs]。安装完成后建议验证安装是否成功python -c import airflow.providers.amazon; print(airflow.providers.amazon.__version__)更详细的安装说明包括从源码安装、升级、卸载参见 Airflow 官方安装文档原文档通过:doc:apache-airflow:installation/index 交叉引用即 Airflow 总仓库的 安装指南。第三步配置 AWS 连接Setup Connection安装完成后还需要在 Airflow 中配置一个 AWS 连接Connection供 Hook/Operator 获取凭证。Amazon Provider 的完整连接说明位于 providers/amazon/docs/connections/aws.rst本节提炼其中最核心的部分。3.1 默认连接 ID 与凭证解析规则默认连接 ID 是aws_default。如果运行 Airflow 的机器上存在${HOME}/.aws/下的凭证文件且该默认连接的 Login/Password 字段为空Airflow 会自动读取其中的凭证。若未运行airflow connections create-default-connections命令大概率没有aws_default连接。此时 Amazon Provider 组件会回退到 boto3 默认凭证策略环境变量、IAM 实例配置文件等。如果需要显式使用该策略应在调用组件时传conn_idNone而不是传一个不存在的连接 ID否则日志中会出现告警。历史变更提醒旧版本安装时aws_default的 extras 字段默认带{region_name: us-east-1}当前版本已不再如此region 需要手动在连接界面设置或通过AWS_DEFAULT_REGION环境变量指定。3.2 连接字段与 Extra 参数在 Airflow 连接界面中或通过 URI / 环境变量AWS 连接的核心字段为AWS Access Key IDLogin初始连接的访问密钥 IDAWS Secret Access KeyPassword初始连接的访问密钥ExtraJSON 字典可配置以下参数全部可选创建初始boto3.session.Session的参数aws_access_key_id、aws_secret_access_key、aws_session_token外部临时凭证需自行续期、region_name、profile_nameAssume Role 相关role_arn指定后通过assume_role_method获取临时安全凭证、assume_role_methodassume_role/assume_role_with_saml/assume_role_with_web_identity默认assume_role、assume_role_kwargsWeb Identity 联邦相关assume_role_with_web_identity_federationfile或google、assume_role_with_web_identity_token_file、assume_role_with_web_identity_federation_audience传给boto3.session.Session.client/resource的参数config_kwargs构造botocore.config.Config例如设置signature_version: unsigned匿名访问公共资源、endpoint_url全局端点如指向 LocalStack/MinIO 等兼容服务、verify是否校验 SSL 证书可传False或 CA bundle 路径按服务细分配置service_config可按s3、sts、ec2等单独指定endpoint_url等参数。3.3 通过代码创建连接以下代码演示如何用 Python 创建 AWS 连接并生成环境变量与 URI源自 connections/aws.rst 示例import os from airflow.models.connection import Connection conn Connection( conn_idsample_aws_connection, conn_typeaws, loginYOUR_AWS_ACCESS_KEY_ID, # AWS Access Key ID passwordYOUR_AWS_SECRET_ACCESS_KEY, # AWS Secret Access Key extra{ region_name: eu-central-1, }, ) env_key fAIRFLOW_CONN_{conn.conn_id.upper()} conn_uri conn.get_uri() print(f{env_key}{conn_uri}) # AIRFLOW_CONN_SAMPLE_AWS_CONNECTIONaws://YOUR_AWS_ACCESS_KEY_ID:YOUR_AWS_SECRET_ACCESS_KEY/?region_nameeu-central-1 os.environ[env_key] conn_uri print(conn.test_connection()) # 校验连接凭证使用 CLI 添加连接时有一个已知坑点当 login/password/host/port 都为空时URI 中需要额外带一个例如airflow connections add aws_conn --conn-uri aws:///?region_nameeu-west-13.4 连接测试的准确含义重要在 Airflow UI / API 中测试 AWS 连接时底层调用的是AwsGenericHook.test_connection见 base_aws.py 实现。其真实逻辑是创建 boto3 session 后调用AWS STS 的GetCallerIdentityAPI根据返回的 HTTP 状态码判断凭证是否有效。由此得出两个重要结论连接测试只能验证凭证是否有效无法验证该凭证是否有权限访问某个具体 AWS 服务如 S3、EC2当使用MinIO、LocalStack 等 AWS API 兼容服务时连接测试失败并不代表凭证错误——这些兼容服务大多只实现了部分 AWS API很多没有实现 STSGetCallerIdentity测试接口自然无法工作。另外从get_ui_field_behaviourbase_aws.py可以看到连接界面会隐藏 host/schema/port 字段并将 Login/Password 重新标注为 “AWS Access Key ID” / “AWS Secret Access Key”同时给出 Extra 字段的 JSON 占位示例含 region、profile、retry、role_arn、endpoint_url 等可作为手工填写的参考模板。3.5 常见认证形态速览IAM 用户密钥对直接在 Login/Password 填入 Access Key ID 与 Secret Access KeyIAM 实例配置文件创建一个空连接aws://或{conn_type: aws}由 boto 默认凭证链自动从实例元数据获取凭证Assume RoleExtra 中指定role_arn例如{ role_arn: arn:aws:iam::112223334444:role/my_role, region_name: ap-southeast-2 }Web Identity 联邦file 模式{ role_arn: arn:aws:iam::112223334444:role/my_role, assume_role_method: assume_role_with_web_identity, assume_role_with_web_identity_federation: file, assume_role_with_web_identity_token_file: /path/to/access_token }SAML 联邦Extra 中配置assume_role_with_saml容器含principal_arn、idp_url、idp_auth_method、mutual_authentication、idp_request_kwargs、idp_request_retry_kwargs、saml_response_xpath等完整示例见 connections/aws.rst。注意该模式依赖requests_gssapi库必要时需pip uninstall python-gssapi pip install gssapipython-gssapi已过时且与 Airflow 使用的paramiko存在版本冲突EKS IRSA在 EKS 上运行时创建一个所有字段为空的 AWS 连接Airflow 即可遵循 boto3 默认凭证链通过 Pod 环境变量AWS_ROLE_ARN、AWS_WEB_IDENTITY_TOKEN_FILE获取 IAM 角色凭证一旦在连接里显式设置了role_arn等字段Airflow 将手动创建 session不再走 boto3 默认流程。底层机制组件如何消费连接与凭证理解前置任务的第三步“配置连接”之后再补充一点源码层面的机理帮助你判断配置是否生效。所有 Amazon Provider 的 Hook 都继承自 base_aws.py 中定义的两个核心类BaseSessionFactory第 104 行起负责 boto3 session 的创建支持同步与异步会话覆盖绝大多数 AWS 认证方式用户也可以通过子类化并重写create_session/_create_basic_session实现自定义联邦认证对应 aws.rst 中的 Session Factory 一节配合[aws] session_factory my_company.aws.MyCustomSessionFactory配置项生效AwsGenericHook第 473 行起通用 AWS Hook 基类提供test_connection、客户端/资源获取、重试装饰器等公共能力AwsBaseHook第 1064 行起AwsGenericHook面向同步 boto3 client/resource 的具体化子类。当某个 Operator如EC2StartInstanceOperator被调度执行时其内部 Hook 会读取连接 ID → 通过BaseSessionFactory依据连接字段Access Key、region、profile、role_arn 等构建 boto3 session → 创建对应服务的 client/resource → 发起 AWS API 调用。因此前置任务中“连接配置”的质量直接决定了后续所有组件调用的成败。常见问题排查清单现象可能原因与处理DAG 运行报NoCredentialsError/ 找不到凭证未配置aws_default连接且环境无默认凭证链补建连接或在调用时显式传conn_idNone走 boto3 默认链连接测试失败但 DAG 运行正常目标服务为 LocalStack/MinIO 等兼容服务未实现 STSGetCallerIdentity属预期行为使用 EKS 却提示凭证无效连接中显式设置了role_arn等字段绕过了 boto3 默认链改用全空连接以启用 IRSA出现 ThrottlingExceptionAWS API 配额限制触发可在连接 Extra 中通过config_kwargs.retries设置{mode: standard, max_attempts: 10}或在~/.aws/config中配置retry_mode standard、max_attempts 10或设置环境变量AWS_RETRY_MODEstandard、AWS_MAX_ATTEMPTS10详见 connections/aws.rst连接界面找不到 Access Key 字段UI 对 aws 类型连接做了字段隐藏与重命名登录/密码即密钥对详见get_ui_field_behaviour总结Amazon Provider 的前置任务高度统一先在 AWS 控制台/CLI 备好资源再pip install apache-airflow[amazon]安装 Provider 及其 API 依赖最后按 AWS 连接文档 配置连接建议从默认连接 IDaws_default开始。这三步完成后即可进入仓库中任意组件文档如 operators 目录、transfer 目录按需编写 DAG。若需深入定制可进一步阅读 base_aws.py 中的 Session Factory 扩展机制与连接测试实现让认证链路完全贴合你的安全与合规要求。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表