尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

使用 Apache Airflow 阿里云 Provider 的 OSS Operators 管理 OSS 存储桶与对象

发布时间:2026/9/13 2:32:43

资讯中心
01
ARTICLE

使用 Apache Airflow 阿里云 Provider 的 OSS Operators 管理 OSS 存储桶与对象

使用 Apache Airflow 阿里云 Provider 的 OSS Operators 管理 OSS 存储桶与对象
使用 Apache Airflow 阿里云 Provider 的 OSS Operators 管理 OSS 存储桶与对象【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文围绕 Apache Airflow 的apache-airflow-providers-alibaba中 Alibaba Cloud OSSObject Storage Service相关的 Operators 与 Sensor 展开说明如何利用 Airflow 以编程方式创建/删除 OSS 存储桶、上传/下载/删除对象并等待对象出现。读完本文你将掌握全部 OSS 相关算子OSSCreateBucketOperator、OSSDeleteBucketOperator、OSSUploadObjectOperator、OSSDownloadObjectOperator、OSSDeleteObjectOperator、OSSDeleteBatchObjectOperator与OSSKeySensor的参数语义、底层 Hook 实现机制以及如何配置oss_default连接完成鉴权并可直接在项目中运行仓库自带的示例 DAG。一、整体概览OSS 集成提供的算子与传感器根据 OSS Operators 文档Airflow 与阿里云对象存储 OSS 的集成提供了以下用于创建与操作 OSS 存储桶的组件全部位于airflow.providers.alibaba.cloud命名空间下组件完整类路径职责OSSKeySensorairflow.providers.alibaba.cloud.sensors.oss_key.OSSKeySensor轮询等待某个 key对象在 OSS 存储桶中出现OSSCreateBucketOperatorairflow.providers.alibaba.cloud.operators.oss.OSSCreateBucketOperator创建 OSS 存储桶OSSDeleteBucketOperatorairflow.providers.alibaba.cloud.operators.oss.OSSDeleteBucketOperator删除 OSS 存储桶OSSUploadObjectOperatorairflow.providers.alibaba.cloud.operators.oss.OSSUploadObjectOperator将本地文件上传为 OSS 对象OSSDownloadObjectOperatorairflow.providers.alibaba.cloud.operators.oss.OSSDownloadObjectOperator将 OSS 对象下载到本地文件OSSDeleteBatchObjectOperatorairflow.providers.alibaba.cloud.operators.oss.OSSDeleteBatchObjectOperator批量删除 OSS 对象OSSDeleteObjectOperatorairflow.providers.alibaba.cloud.operators.oss.OSSDeleteObjectOperator删除单个 OSS 对象这些算子全部继承自airflow.providers.common.compat.sdk.BaseOperator而OSSKeySensor继承自BaseSensorOperator其底层统一委托给 OSSHook 完成实际的 OSS API 调用见 operators/oss.py。二、创建与删除 OSS 存储桶2.1 目的文档中给出的是一个完整的示例 DAG使用OSSCreateBucketOperator以指定的存储桶名称创建新桶随后使用OSSDeleteBucketOperator将其删除。该示例同时演示了如何把创建桶与删除桶两个任务串成依赖链。2.2 任务定义示例仓库中的系统测试 DAG example_oss_bucket.py 完整实现了这一流程from datetime import datetime from airflow.models.dag import DAG from airflow.providers.alibaba.cloud.operators.oss import OSSCreateBucketOperator, OSSDeleteBucketOperator REGION os.environ.get(REGION, default_region) with DAG( dag_idoss_bucket_dag, start_datedatetime(2021, 1, 1), scheduleNone, default_args{bucket_name: your bucket, region: your region}, max_active_runs1, tags[example], catchupFalse, ) as dag: create_bucket OSSCreateBucketOperator(task_idtask1, regionREGION) delete_bucket OSSDeleteBucketOperator(task_idtask2, regionREGION) create_bucket delete_bucket要点解读region是必填参数用于指定创建/删除桶所在的地域示例中通过环境变量REGION注入。bucket_name为可选参数默认None。示例把bucket_name放在 DAG 的default_args中说明算子会优先使用任务显式传入的值当未显式提供时Hook 层会自动回退到连接配置中的默认桶名详见下文第四节桶名解析。oss_conn_id默认值为oss_default即阿里云连接文档中声明的默认连接 ID。任务依赖create_bucket delete_bucket保证先建后删这是典型的资源生命周期管理 DAG 写法。2.3 源码级实现从 operators/oss.py 可以看到两个算子的execute都只是简单地把参数转发给 Hookdef execute(self, context: Context): oss_hook OSSHook(oss_conn_idself.oss_conn_id, regionself.region) oss_hook.create_bucket(bucket_nameself.bucket_name)而OSSHook.create_bucket/delete_bucket最终构造 OSS SDK v2 的PutBucketRequest/DeleteBucketRequest并调用客户端见 hooks/oss.py异常统一包装为AirflowException抛出保证失败能被 Airflow 正确捕获并标记任务失败。单元测试 test_oss.py 也验证了这一点算子execute时会以oss_conn_id和region构造OSSHook并调用create_bucket(bucket_name...)/delete_bucket(bucket_name...)。三、对象级操作上传、下载、删除除存储桶生命周期管理外OSS 集成还提供了对象Object级的四个算子。仓库中的系统测试 DAG example_oss_object.py 演示了完整用法create_object OSSUploadObjectOperator( fileyour local file, keyyour oss key, task_idtask1, regionREGION, ) download_object OSSDownloadObjectOperator( fileyour local file, keyyour oss key, task_idtask2, regionREGION, ) delete_object OSSDeleteObjectOperator( keyyour oss key, task_idtask3, regionREGION, ) delete_batch_object OSSDeleteBatchObjectOperator( keys[obj1, obj2, obj3], task_idtask4, regionREGION, ) create_object download_object delete_object delete_batch_object各算子参数语义如下依据 operators/oss.pyOSSUploadObjectOperator(key, file, region, bucket_nameNone, oss_conn_idoss_default)key是对象在 OSS 中的路径file是本地待上传文件路径。执行时调用 Hook 的upload_local_file底层使用 OSS SDK v2 的uploader().upload_file()支持断点续传等高级能力见 hooks/oss.py。OSSDownloadObjectOperator(key, file, region, bucket_nameNone, oss_conn_idoss_default)key为要下载的对象键file为本地保存路径路径文件名。执行时调用download_file底层使用downloader().download_file()成功时返回本地文件名失败时记录错误并返回None见 hooks/oss.py。OSSDeleteObjectOperator(key, region, bucket_nameNone, oss_conn_idoss_default)删除单个对象对应DeleteObjectRequest见 hooks/oss.py。OSSDeleteBatchObjectOperator(keys, region, bucket_nameNone, oss_conn_idoss_default)keys是对象键列表一次删除多个对象对应DeleteMultipleObjectsRequest见 hooks/oss.py。单元测试 test_oss.py 分别验证了上述算子对 Hook 方法的调用关系可作为编写自定义 DAG 时的参考。四、OSSKeySensor等待对象出现OSSKeySensor用于等待某个 key对象出现在 OSS 存储桶中典型的应用场景是上游任务产生文件到 OSS下游任务消费该文件之间的衔接。4.1 参数说明依据 oss_key.pybucket_key必填等待的 key。支持两种写法完整的oss://风格 URL此时必须将bucket_name留为None传感器会自动从 URL 解析出桶名和 key从桶根目录开始的相对路径此时必须显式提供bucket_name。region必填OSS 地域。bucket_name可选桶名。oss_conn_id默认oss_default。继承自BaseSensorOperator的轮询参数如poke_interval、timeout、mode等均可正常使用。4.2 关键行为从 oss_key.py 的实现可以看到poke中的参数校验逻辑若未提供bucket_name而bucket_key又不是完整的oss://URLnetloc为空直接抛出AirflowException提示请提供 bucket_name若提供了bucket_name却又传入完整oss://URL同样抛出异常提示bucket_key 应为相对路径而非完整 URL校验通过后传感器记录日志Poking for key : oss://bucket/key并调用self.hook.object_exists(...)判断对象是否存在。object_exists在 Hook 中通过is_object_exist实现见 hooks/oss.py。另外两点值得注意template_fields (bucket_key, bucket_name)说明bucket_key与bucket_name支持 Jinja 模板渲染可以在传感器中引用{{ ds }}等上下文变量动态生成路径。get_hook属性已标记为deprecated抛出AirflowProviderDeprecationWarning官方推荐直接使用hook属性见 oss_key.py。五、底层 HookOSSHook 与连接配置5.1 连接配置Connection所有 OSS 算子/传感器共用oss_default连接默认连接 ID配置方式见 Alibaba Cloud Connection 文档Connection TypeossSchema可选指定 OSS Hook 使用的默认存储桶名称。当算子在bucket_name参数缺省时Hook 会从连接的 Schema 自动取桶名。Extra可选JSON 字典鉴权与地域配置当前仅支持AKAccessKey方式{ auth_type: AK, access_key_id: your access key id, access_key_secret: your access key secret }Extra 支持以下参数参数说明是否必填auth_type访问阿里云资源的认证类型当前仅支持AK必填access_key_id阿里云用户的 AccessKey ID必填access_key_secret阿里云用户的 AccessKey Secret必填endpoint自定义 OSS Endpoint可选缺省时自动推导为oss-region.aliyuncs.com可选region默认地域缺省时由算子的region参数提供可选5.2 Hook 的鉴权与客户端构造从 hooks/oss.py 可以看出客户端构造逻辑使用 OSS SDK v2alibabacloud_oss_v2的oss.config.load_default()加载默认配置config.region取自 Hook 的region属性——其优先级是算子显式传入的region 连接 Extra 中的regionget_default_region()见 hooks/oss.pyconfig.endpoint默认拼接为oss-region.aliyuncs.com也可在连接 Extra 中用endpoint覆盖凭据通过get_credential()生成StaticCredentialsProvider其中校验了auth_type必须为AK并要求access_key_id与access_key_secret均非空见 hooks/oss.py。5.3 桶名解析装饰器Hook 层通过两个装饰器实现了优雅的缺省回退见 hooks/oss.pyprovide_bucket_name当bucket_name参数为None且连接存在时从连接的schema字段取出默认桶名注入调用。unify_bucket_name_and_key当bucket_name为None时调用静态方法parse_oss_url把oss://bucket/key形式的 URL 拆解为桶名与 key使用urlsplit解析若 URL 缺少 netloc 则抛出AirflowException见 hooks/oss.py。也就是说Hook 的多数方法既支持bucket_name key分离传参也支持把完整oss://URL 直接作为key传入这一灵活性向上透传给了各算子。六、在项目中落地从安装到运行示例 DAG6.1 安装 Provider本仓库中阿里云 Provider 的源码位于 providers/alibaba实际发布包名为apache-airflow-providers-alibaba。在 Airflow 环境中可通过 pip 安装该 Providerpip install apache-airflow-providers-alibaba从源码安装的方式可参考 installing-providers-from-sources.rst。6.2 配置连接在 Airflow UI 的 Admin → Connections 中新增一条连接Conn Idoss_defaultConn TypeossSchema默认桶名可选Extra{ auth_type: AK, access_key_id: 你的 AccessKey ID, access_key_secret: 你的 AccessKey Secret, region: cn-hangzhou }6.3 运行示例 DAG仓库提供了两个可直接参考的系统测试 DAGexample_oss_bucket.py创建桶 → 删除桶example_oss_object.py上传对象 → 下载对象 → 删除对象 → 批量删除对象。将示例 DAG 中的your bucket、your local file、your oss key等占位符替换为真实值并把REGION环境变量设置为你的地域即可在 Airflow 中调度运行。运行系统测试的完整方法可参考 system_tests.rst。七、小结Apache Airflow 阿里云 Provider 的 OSS 集成以OSSHook为统一底座向上封装出 6 个算子与 1 个传感器覆盖了对象存储最常见的生命周期操作建桶、删桶、上传、下载、单删、批量删除与对象就绪等待。理解三个关键点即可举一反三参数回退链bucket_name缺省时 → 取连接 Schemaregion缺省时 → 取连接 Extra 的region均未配置则报错。URL 与分离传参二选一传完整oss://URL 时必须留空bucket_name反之亦然传感器与 Hook 都会做严格校验。鉴权方式当前仅支持 AK 静态凭据auth_type: AK凭据统一存放在oss_default连接的 Extra 字段中endpoint默认按地域推导、可自定义覆盖。配合仓库中的 算子源码、Hook 实现 与 单元测试你可以在此基础上快速构建自己的 OSS 数据流转工作流。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。