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

资讯详情

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

Alibaba Cloud 连接配置完全指南:Apache Airflow Alibaba Provider 的认证机制与 Connection 配置实战

Alibaba Cloud 连接配置完全指南:Apache Airflow Alibaba Provider 的认证机制与 Connection 配置实战 Alibaba Cloud 连接配置完全指南Apache Airflow Alibaba Provider 的认证机制与 Connection 配置实战【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的apache-airflow-providers-alibaba提供方Provider封装了阿里云 OSS、MaxCompute、AnalyticDB for MySQL Spark 等服务的 Hook、Operator 与 Sensor。本文基于提供方官方连接文档结合仓库源码完整讲解 Alibaba Cloud Connection 的认证机制、默认连接 ID、Schema 与 Extra 字段配置规范并给出可在 Airflow 中直接落地的配置步骤与源码级原理佐证。1. 认证机制基于 STS 与签名 URL 的阿里云鉴权根据 alibaba.rst 连接文档向阿里云资源发起访问时的认证可以基于安全令牌服务STSSecurity Token Service或签名 URLsigned URL完成。STS 是阿里云提供的临时凭证颁发服务适合在 DAG 运行期间动态获取短期有效的访问凭证而签名 URL 则是针对 OSS 等对象存储服务的临时授权方式适用于分享或限时访问场景。从源码结构来看当前提供方实际落地的是基于 AccessKey 的静态凭证认证auth_type AK在 OSSHook 中get_credential()方法从连接的extra_dejson中读取auth_type、access_key_id、access_key_secret三个字段校验通过后构造oss.credentials.StaticCredentialsProvider交给 OSS SDK v2 客户端使用AnalyticDBSparkHook 的get_adb_spark_client()同样校验这三个字段并用其构造alibabacloud_tea_openapi.models.Config初始化 ADB Spark REST API 客户端。需要特别指出的是尽管文档声明认证“可能”基于 STS 或签名 URL但当前版本源码中的校验逻辑只放行AK一种类型——auth_type缺失或取值非AK时两个 Hook 都会抛出ValueError(Unsupported auth_type: ...)。因此在实际配置中auth_type必须显式设置为AK并将阿里云账号的 AccessKey ID 与 AccessKey Secret 填入对应字段。2. 默认连接 IDoss_default 与 adb_spark_default文档明确给出了两个默认连接 IDoss_default供 OSS Hook / Operator / Sensor 使用adb_spark_default供 AnalyticDB for MySQL Spark Hook / Operator 使用。这一约定在源码中得到了一一印证类默认连接 ID连接类型conn_type源码位置OSSHookoss_defaultosshooks/oss.pyAnalyticDBSparkHookadb_spark_defaultadb_sparkhooks/analyticdb_spark.py此外提供方还内置了另外两个连接类型可视为对文档内容的补充AlibabaBaseHook 的默认连接 ID 为alibabacloud_default、连接类型为alibaba_cloudMaxCompute 的 MaxComputeHook 默认连接 ID 为maxcompute_default、连接类型为maxcompute。这些默认值均通过类属性default_conn_name暴露给 Airflow 的 Connection 管理机制意味着在配置连接时如果不显式指定 ID将自动回落到对应的默认值。在 UI 或代码中创建连接时可以自定义连接 ID例如my_oss_conn然后在 Operator/Hook 构造时通过oss_conn_id、adb_spark_conn_id参数显式传入覆盖默认值。3. Schema 字段为 OSS Hook 指定默认 Bucket文档指出连接的Schema可选字段用于指定 OSS Hook 使用的默认 Bucket 名称。从源码看这一行为由装饰器provide_bucket_name实现见 hooks/oss.py当调用OSSHook中标注了该装饰器的方法如object_exists、get_bucket、load_string、upload_local_file、download_file、create_bucket等且未显式传入bucket_name时装饰器会读取连接的schema字段并自动填充为bucket_name参数。也就是说若在连接配置中将 Schema 设为my-data-bucket那么在使用OSSHook时未显式指定bucket_name的操作都会自动作用于该 Bucket从而减少 DAG 中的重复参数。同时OSSHook还提供unify_bucket_name_and_key装饰器见 hooks/oss.py当只传了形如oss://bucket-name/path/to/key的完整 URL 而未传bucket_name时会通过parse_oss_url()自动拆解出 Bucket 与 Key。这一机制让 OSS 操作既支持“显式指定 bucket_name key”的写法也支持“单个 oss:// 完整路径”的写法两种风格在 DAG 中皆可混用。4. Extra 字段JSON 格式的认证与连接参数文档规定连接的Extra可选字段用于存放 Alibaba Cloud Connection 的附加参数以JSON 字典形式填写且下列参数均为可选auth_type访问阿里云资源使用的认证类型当前仅支持AKaccess_key_id阿里云用户的 Access Key IDaccess_key_secret阿里云用户的 Access Key Secret。文档给出的 Extra 字段标准示例为{ auth_type: AK, access_key_id: , access_key_secret: }4.1 源码层面的字段校验与扩展字段虽然文档只列出三个参数但从源码可以确认extra中还支持若干扩展字段且校验逻辑非常严格region区域OSSHook.get_default_region()与AnalyticDBSparkHook.get_default_region()都会从extra中读取region未配置时抛出ValueError。OSS 客户端在 hooks/oss.py 中按foss-{self.region}.aliyuncs.com自动推导 endpointADB Spark 客户端则按fadb.{self.region}.aliyuncs.com构造 endpoint。endpointOSS 专用OSSHook._get_client()优先读取extra中的endpoint若未提供才回落到按 region 推导的默认 endpoint。因此使用 OSS 自定义域名、内网 endpoint 或专有网络VPCendpoint 时可通过该字段显式覆盖。project / endpointMaxCompute 专用MaxComputeHook 通过fallback_to_default_project_endpoint装饰器支持从extra中读取project与endpoint作为调用get_client、run_sql、get_instance、stop_instance时的默认值extra中还支持access_key_id、access_key_secret两个凭证字段与文档一致。此外OSS extra中缺失access_key_id或access_key_secret时get_credential() 会抛出包含连接 ID 的ValueError便于快速定位是哪条连接配置不完整。4.2 UI 表单化支持为了降低手动拼 JSON 的出错概率提供方在 Airflow 连接编辑表单中内置了可视化字段AlibabaBaseHook.get_connection_form_widgets()见 base_alibaba.py为access_key_id、access_key_secret渲染了密码型输入框BS3PasswordFieldWidget凭证不会明文展示MaxComputeHook在此基础上追加了project、endpoint两个文本框并通过get_ui_field_behaviour()将host、schema、login、password、port、extra等通用字段隐藏引导用户只填写必要项。也就是说在 Airflow Web UI 的 Connections 页面选择对应连接类型后可以直接在表单中填写凭证最终同样会以 JSON 形式写入 Extra。5. 实战配置步骤与验证5.1 通过 Airflow CLI 创建 OSS 连接若使用 Airflow CLI可以执行airflow connections add oss_default \ --conn-type oss \ --conn-schema my-data-bucket \ --conn-extra {auth_type: AK, access_key_id: LTAI5t..., access_key_secret: your-secret, region: cn-hangzhou}要点--conn-schema对应文档中的 Schema 字段作为 OSS 默认 Bucket--conn-extra传入 JSON除文档要求的三个参数外务必补充region否则get_default_region()会抛错如需自定义 endpoint可在 JSON 中追加endpoint: oss-cn-hangzhou-internal.aliyuncs.com等值。5.2 通过 Airflow CLI 创建 AnalyticDB Spark 连接airflow connections add adb_spark_default \ --conn-type adb_spark \ --conn-extra {auth_type: AK, access_key_id: LTAI5t..., access_key_secret: your-secret, region: cn-hangzhou}AnalyticDBSparkHook.get_default_region()同样强制要求region因此该字段不可省略。5.3 在 DAG 中使用连接连接配置完成后即可在 DAG 中通过默认连接 ID 直接使用from airflow.providers.alibaba.cloud.hooks.oss import OSSHook from airflow.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkBatchOperator # OSS 读写自动使用 oss_default 连接 Schema 中的默认 bucket hook OSSHook(regioncn-hangzhou) hook.load_string(keylogs/hello.txt, contenthello alibaba) # 提交 AnalyticDB Spark 批任务 spark_pi AnalyticDBSparkBatchOperator( task_idspark_pi, filelocal:///tmp/spark-examples.jar, class_nameorg.apache.spark.examples.SparkPi, cluster_idyour cluster id, rg_nameyour resource group name, )仓库自带的系统测试 DAG example_adb_spark_batch.py 展示了更完整的批量提交写法两个 Spark 示例任务SparkPi、SparkLR串行执行cluster_id、rg_name、region均可放入default_args统一管理。5.4 错误排查提示结合源码校验逻辑配置完成后若运行时报错可按下表快速定位报错信息节选原因No auth_type specified in extra_config.Extra 中缺少auth_typeUnsupported auth_type: ...auth_type取值不是AKNo access_key_id is specified for connection: ...Extra 中缺少access_key_idNo access_key_secret is specified for connection: ...Extra 中缺少access_key_secretNo region is specified for connection: ...Extra 中缺少region6. 关联资源速查连接文档原文providers/alibaba/docs/connections/alibaba.rstOSS Hook 实现与装饰器cloud/hooks/oss.pyAnalyticDB Spark Hook 与提交参数构造cloud/hooks/analyticdb_spark.py通用认证基类与 UI 表单cloud/hooks/base_alibaba.pyMaxCompute Hookproject/endpoint 扩展用法cloud/hooks/maxcompute.pyOSS 键存在性 Sensorcloud/sensors/oss_key.py系统测试示例 DAGtests/system/alibaba/example_adb_spark_batch.py7. 安装与使用前提使用该连接前需先安装提供方包在已安装 Apache Airflow 的环境上pip install apache-airflow-providers-alibaba根据 providers/alibaba/docs/index.rst该提供方要求 Apache Airflow 2.11.0并依赖apache-airflow-providers-common-compat1.13.0、alibabacloud-oss-v21.2.0、alibabacloud_adb202112011.0.0、alibabacloud_tea_openapi0.3.7以及pyodps等运行时库。若使用 OSS 之外的组件如 MaxCompute请确认对应 Python 依赖已一并安装。本文所述配置方式均以当前仓库中的提供方实现为准连接文档只声明了auth_type/access_key_id/access_key_secret三个 Extra 参数而实际落地时region等字段同样必不可少这是仓库源码校验逻辑所决定的行为配置时务必留意。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表