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/airflow
Apache 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_default@my-bucket/")URI 中的用户名部分(aws_default)代表 Airflow 的connection id,是可选的;也可以换成独立的conn_id关键字参数,两种写法完全等价:
# Equivalent to the previous example. base = ObjectStoragePath("s3://my-bucket/", conn_id="aws_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, b"data") 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(protocol="dbfs", fs=DBFSFileSystem(instance="myinstance", token="mytoken")) 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 Pathlib(upath)之上,因此大部分操作与你操作本地文件系统的方式一致。本节只列出与标准 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-L206 |
checksum | 返回文件的校验和 | path.py#L273-L276 |
container | bucket的别名 | path.py#L198-L200 |
fs | 便捷属性,返回已实例化的(Airflow 认证后的)文件系统 | path.py#L175-L178 |
key | 返回对象 key(按惯例去掉前导斜杠,保留尾部斜杠以支持目录语义) | path.py#L208-L214 |
namespace | 返回对象的命名空间,通常是协议加桶名,如s3://bucket | path.py#L216-L218 |
path | 供文件系统实例使用的 fsspec 兼容路径 | 继承自 UPath |
protocol | fsspec 协议名 | 继承自 UPath |
read_block | 从文件指定偏移读取字节块 | path.py#L278-L317 |
sign | 生成代表该路径的签名 URL,用于委托凭证(支持临时 URL 的实现可用) | path.py#L319-L336 |
size | 返回文件字节大小 | path.py#L338-L340 |
storage_options | 实例化底层文件系统所用的存储选项 | 继承自 UPath |
ukey | 文件属性的哈希,用于判断文件是否变化 | path.py#L269-L271 |
其中read_block(offset, length, delimiter=None)的行为值得展开:从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(..., recursive=True)展开目录树,跳过空目录,逐文件调用_cp_file; - 远端到远端的复制,目标 key 与源保持一致——即
s3://src_bucket/foo/bar会复制到gcs://dst_bucket/foo/bar,而非gcs://dst_bucket/bar; - 目标若已存在同名文件或目录,按 fsspec 语义覆盖。
move(path.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_id="aws_default") conn = duckdb.connect(database=":memory:") conn.register_filesystem(path.fs) conn.execute(f"CREATE 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_default@airflow-tutorial-data/")get_air_quality_data任务调用公开 API 拉取空气质量数据,先base.mkdir(exist_ok=True)确保桶/目录存在,再用base / f"air_quality_{formatted_date}.parquet"拼接按日期命名的路径,以二进制写模式path.open("wb")写入 parquet,最后把路径作为返回值交给下游(示例 L66-L103):
@task def get_air_quality_data(logical_date=None) -> ObjectStoragePath: ... # ensure the bucket exists base.mkdir(exist_ok=True) formatted_date = logical_date.format("YYYYMMDD") path = base / f"air_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.ProxyUPath(path.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 == "s3"(test_path.py#L37-L54)。
血缘(Lineage)自动采集
open()返回的是一个_TrackingFileWrapper(path.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),仅供参考