ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Airflow Object Storage 抽象:用 ObjectStoragePath 统一操作 S3、GCS 与 Azure Blob

Apache Airflow Object Storage 抽象:用 ObjectStoragePath 统一操作 S3、GCS 与 Azure Blob Apache Airflow Object Storage 抽象用 ObjectStoragePath 统一操作 S3、GCS 与 Azure Blob【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 从 2.8.0 起提供了内置的对象存储抽象层将 S3、GCS、Azure Blob 等云对象存储统一封装为ObjectStoragePath路径对象让 DAG 开发者可以用一套近乎 pathlib 的 API 读写各种对象存储无需为不同云厂商编写分支代码。本文以官方文档 objectstorage.rst 为骨架结合仓库中ObjectStoragePath的实现源码与示例 DAG完整讲解该抽象的用法、配置方式、Path API、扩展操作、跨存储复制移动以及外部集成方案读完即可在 DAG 中直接落地使用。对象存储的本质它不是真正的文件系统对象存储是云厂商提供的主流持久化存储形态但它并非经典意义上的 POSIX 文件系统。为了在海量数据数百 PB 级下消除单点故障对象存储用更简单的object-name data映射模型取代了传统文件系统的目录树为了支持远程访问对对象的操作通常以相对较慢的HTTP REST 请求形式提供。Airflow 在 S3、GCS、Azure Blob 等对象存储之上提供了一层通用抽象目标是在 DAG 中使用多种对象存储系统时无需修改业务代码可以配合shutil等大多数标准 Python 模块使用它们能操作 file-like 对象。由于对象存储不是真正的文件系统使用时有几个与本地文件系统明显不同的关键点设计 DAG 时必须留意没有原子的重命名操作移动文件实际是先复制、再删除。如果复制失败源文件可能丢失目录是模拟出来的例如列出一个目录可能需要在桶内列出全部对象再按前缀过滤速度可能很慢文件内 seek定位可能需要较高的调用开销、影响性能甚至可能根本不支持。Airflow 依赖 fsspec 的导入与类定义。但从源码结构看缓存与性能优化只是辅助设计 DAG 时仍需正视对象存储的上述固有限制。基本使用从一条 URI 开始使用对象存储的第一步是用目标对象的 URI 实例化一个ObjectStoragePath。例如指向 S3 中的某个桶from airflow.sdk import ObjectStoragePath base ObjectStoragePath(s3://aws_defaultmy-bucket/)URI 中的用户名部分aws_default代表 Airflow 的connection id是可选的也可以换成独立的conn_id关键字参数两种写法完全等价# Equivalent to the previous example. base ObjectStoragePath(s3://my-bucket/, conn_idaws_default)从源码看ObjectStoragePath.__init__会先用urlsplit解析 URI 中的 userinfo前的部分作为默认 conn_id再用显式传入的conn_id覆盖之随后把 conn_id 从 storage_options 中剥离避免把它误传给不认识的 fsspec 底层文件系统见 task-sdk/src/airflow/sdk/io/path.py#L92-L120。列出文件对象task def list_files() - list[ObjectStoragePath]: files [f for f in base.iterdir() if f.is_file()] return files在目录树中导航/运算符与 pathlib 行为一致用于拼接子路径base ObjectStoragePath(s3://my-bucket/) subdir base / subdir # prints ObjectStoragePath(s3://my-bucket/subdir) print(subdir)打开文件open()返回 file-like 对象可以像本地文件一样读写task def read_file(path: ObjectStoragePath) - str: with path.open() as f: return f.read()通过 XCOM 在任务间传递路径对象存储路径天然适合作为 XCOM 的载体任务产出路径、下游任务消费路径形成清晰的数据流task def create(path: ObjectStoragePath) - ObjectStoragePath: return path / new_file.txt task def write_file(path: ObjectStoragePath, content: str): with path.open(wb) as f: f.write(content) new_file create(base) write write_file(new_file, bdata) read write这一用法有实现层面的支撑ObjectStoragePath实现了serialize/deserialize两个方法把path、conn_id和storage_options打包成可跨任务传输的字典并带版本号控制见 task-sdk/src/airflow/sdk/io/path.py#L484-L500这正是它能安全穿越 XCOM 序列化/反序列化流程的原因。配置连接机制与替代后端连接配置自动下推基本使用场景下对象存储抽象几乎不需要额外配置它完全依赖 Airflow 标准的Connection机制通过conn_id指定要使用的连接连接上的任何设置都会被下推到底层实现。例如使用 S3 时可以在 Connection 中配置aws_access_key_id、aws_secret_access_key还可以通过 extra 传入endpoint_url等参数指定自定义端点。不同对象存储 scheme 的支持取决于你安装的 provider内置开箱即用支持filescheme安装了apache-airflow-providers-google即可使用gcsscheme支持s3需要安装apache-airflow-providers-amazon[s3fs]——因为它依赖aiobotocore而aiobotocore默认不随 botocore 一起安装以免造成依赖冲突。为协议挂载替代后端attach可以为某个 scheme / 协议配置替代后端把backend一个 fsspec 文件系统实例通过attach挂到协议上。例如为dbfsscheme 启用 Databricks 后端from airflow.sdk import ObjectStoragePath from airflow.sdk.io import attach from fsspec.implementations.dbfs import DBFSFileSystem attach(protocoldbfs, fsDBFSFileSystem(instancemyinstance, tokenmytoken)) base ObjectStoragePath(dbfs://my-location/)注意要让后端注册在多个任务间复用必须在DAG 的顶层top-level调用attach否则该后端在其他任务中不可用。从源码看attach的实际行为是以protocol-conn_id无 conn_id 时仅用 protocol为别名建立全局缓存_STORE_CACHE同一别名重复调用直接返回已注册的ObjectStore避免重复创建文件系统实例见 task-sdk/src/airflow/sdk/io/store.py#L130-L163。ObjectStoragePath.fs属性正是通过attach(self.protocol or file, self.conn_id).fs拿到经过 Airflow 连接认证的文件系统见 task-sdk/src/airflow/sdk/io/path.py#L175-L178。Path API与 pathlib 对齐的标准操作对象存储抽象以Path API实现构建在Universal Pathlibupath之上因此大部分操作与你操作本地文件系统的方式一致。本节只列出与标准 Path API 存在差异的操作其余细节可查阅ObjectStoragePath类文档对应实现见 task-sdk/src/airflow/sdk/io/path.py。mkdir在指定路径或桶/容器内创建目录条目。对于没有真正目录概念的系统可能仅为当前实例创建目录条目、并不影响真实文件系统。若parents为True缺失的父级路径会被一并创建。touch在给定路径创建文件或更新时间戳。truncate默认为True即会截断文件若文件已存在当exists_ok为真时操作成功并把修改时间更新为当前时间否则抛出FileExistsError。stat返回一个类似stat_result的对象支持st_size、st_mtime、st_mode等属性同时行为上又像一个字典可提供对象的附加元数据。例如 S3 下会额外返回[ETag, ContentType]等键。如果代码需要跨对象存储移植不要依赖这些扩展元数据。实现上stat会把底层fs.stat结果包装成带protocol、conn_id等信息的stat_result见 task-sdk/src/airflow/sdk/io/path.py#L225-L231并据此实现samefile判断同文件判断。扩展操作超越标准 Path API 的能力以下操作不属于标准 Path API但由对象存储抽象额外支持均可在ObjectStoragePath上直接调用操作说明源码位置bucket返回桶名path.py#L202-L206checksum返回文件的校验和path.py#L273-L276containerbucket的别名path.py#L198-L200fs便捷属性返回已实例化的Airflow 认证后的文件系统path.py#L175-L178key返回对象 key按惯例去掉前导斜杠保留尾部斜杠以支持目录语义path.py#L208-L214namespace返回对象的命名空间通常是协议加桶名如s3://bucketpath.py#L216-L218path供文件系统实例使用的 fsspec 兼容路径继承自 UPathprotocolfsspec 协议名继承自 UPathread_block从文件指定偏移读取字节块path.py#L278-L317sign生成代表该路径的签名 URL用于委托凭证支持临时 URL 的实现可用path.py#L319-L336size返回文件字节大小path.py#L338-L340storage_options实例化底层文件系统所用的存储选项继承自 UPathukey文件属性的哈希用于判断文件是否变化path.py#L269-L271其中read_block(offset, length, delimiterNone)的行为值得展开从offset处开始读取length字节若指定delimiter会确保读写落在紧随offset与offset length之后的分隔符边界上若offset为 0 则从 0 开始返回的字节串包含结尾分隔符若offset length超出 EOF则一直读到 EOF。其 docstring 给出了一个直观例子CSV 文件按换行符切块读取见 path.py#L297-L311。复制与移动跨对象存储的数据搬运copy与move用于把文件或目录从source复制/移动到target其预期行为与 fsspec 规范一致。跨对象存储例如 file - s3复制目录时Airflow 需要遍历目录树、逐个文件处理把每个文件从源流式传输到目标。源码中的_cp_file正是用with self.open(rb) as f1, dst.open(wb) as f2:配合shutil.copyfileobj完成流式拷贝见 task-sdk/src/airflow/sdk/io/path.py#L342-L354。copy的完整分支逻辑如下path.py#L356-L425同存储内直接调用底层fs.copy通常是服务端优化路径本地 - 远端 / 远端 - 本地分别走fs.put/fs.get优化路径远端目录 - 远端目录用fs.expand_path(..., recursiveTrue)展开目录树跳过空目录逐文件调用_cp_file远端到远端的复制目标 key 与源保持一致——即s3://src_bucket/foo/bar会复制到gcs://dst_bucket/foo/bar而非gcs://dst_bucket/bar目标若已存在同名文件或目录按 fsspec 语义覆盖。movepath.py#L443-L465在同存储内直接调用fs.move跨存储时退化为先copy再unlink——这与前文对象存储没有原子重命名的局限一致。此外还提供了copy_into/move_into把文件复制/移动到某个目录内部目标必须是目录否则抛NotADirectoryError。外部集成把 Airflow 的连接能力带给其他工具DuckDB、Apache Iceberg 等许多项目都可以消费这个对象存储抽象通常的接入方式是传入底层的 fsspec 实现。为此ObjectStoragePath暴露了fs属性。例如下面的代码让 DuckDB 复用 Airflow 中配置的连接去连接 S3并直接读取一个由ObjectStoragePath指向的 parquet 文件import duckdb from airflow.sdk import ObjectStoragePath path ObjectStoragePath(s3://my-bucket/my-table.parquet, conn_idaws_default) conn duckdb.connect(database:memory:) conn.register_filesystem(path.fs) conn.execute(fCREATE OR REPLACE TABLE my_table AS SELECT * FROM read_parquet({path});)这段代码之所以成立是因为fs属性返回的正是经由attach(protocol, conn_id)创建的、携带 Airflow 连接凭证的文件系统path.py#L175-L178外部工具注册该文件系统后即获得同样的鉴权与寻址能力无需重复配置凭据。实战示例仓库自带的 Object Storage 教程 DAG仓库在 airflow-core/src/airflow/example_dags/tutorial_objectstorage.py 提供了一个完整的参考 DAG把上述 API 串成了真实的数据流模块顶层创建基础路径conn_id直接内嵌在 URI 中base ObjectStoragePath(s3://aws_defaultairflow-tutorial-data/)get_air_quality_data任务调用公开 API 拉取空气质量数据先base.mkdir(exist_okTrue)确保桶/目录存在再用base / fair_quality_{formatted_date}.parquet拼接按日期命名的路径以二进制写模式path.open(wb)写入 parquet最后把路径作为返回值交给下游示例 L66-L103task def get_air_quality_data(logical_dateNone) - ObjectStoragePath: ... # ensure the bucket exists base.mkdir(exist_okTrue) formatted_date logical_date.format(YYYYMMDD) path base / fair_quality_{formatted_date}.parquet with path.open(wb) as file: df.to_parquet(file) return pathanalyze任务接收上游路径注册path.fs到 DuckDB 并执行 SQL 查询示例 L108-L129完整演示了写入对象存储 - 通过 XCOM 传递路径 - 外部引擎消费的闭环。深入实现ObjectStoragePath 源码剖析继承与认证文件系统注入ObjectStoragePath继承自upath.extensions.ProxyUPathpath.py#L82所有标准路径操作exists、mkdir、iterdir、glob、walk、rename、read_bytes、write_bytes等都委托给内部的__wrapped__UPath。构造时如果指定了 conn_id实现会把 Airflow 认证后的文件系统直接注入__wrapped__._fs_cached从而让所有委托操作一次修复、处处生效而不是为每个方法单独覆写注入失败仅记录 DEBUG 日志错误会在首次使用路径时才暴露path.py#L148-L168。conn_id 的完整传播conn_id会通过_from_upath从父实例传播到所有派生路径/、joinpath、parent、parents、with_name、with_suffix、with_stem等这一行为在单元测试中有系统性覆盖见 task-sdk/tests/task_sdk/io/test_path.py#L63-L103。测试同时验证了 URI 解析规则ObjectStoragePath(s3://bucket/key/part1/part2)会得到bucket bucket、key key/part1/part2、protocol s3test_path.py#L37-L54。血缘Lineage自动采集open()返回的是一个_TrackingFileWrapperpath.py#L40-L79它会拦截 file-like 对象上的read/write调用自动向血缘收集器登记输入/输出资产copy、move在同存储或涉及本地文件时也会显式登记血缘。这意味着在使用对象存储 API 时数据血缘追踪几乎是免费的。向后兼容层Airflow 核心包在 airflow-core/src/airflow/io/init.py 中通过add_deprecated_classes把airflow.io.path.ObjectStoragePath、airflow.io.attach等旧入口映射到新的airflow.sdk位置保证老代码在迁移期间仍可导入。因此新代码应直接使用from airflow.sdk import ObjectStoragePath与文档和示例保持一致。小结Airflow 的对象存储抽象把对象存储不是文件系统这一现实封装成了友好的 Path API用一条带 conn_id 的 URI 实例化ObjectStoragePath即可完成列目录、导航、读写、复制移动、签名 URL、按块读取等操作并把 Airflow 的连接配置自动下推给底层 fsspec 文件系统通过fs属性还能把同一套凭证能力借给 DuckDB、Apache Iceberg 等外部引擎。官方文档 objectstorage.rst、示例 DAG tutorial_objectstorage.py 以及实现源码 task-sdk/src/airflow/sdk/io/path.py 共同构成了从上手使用到原理理解的完整学习路径。需要特别提醒的是对象存储没有原子重命名、目录操作可能昂贵、seek 开销不可忽视——设计 DAG 时始终把这些局限记在心里才能写出真正健壮的数据管道。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表