news 2026/9/10 12:29:55

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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/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两个方法,把pathconn_idstorage_options打包成可跨任务传输的字典,并带版本号控制(见 task-sdk/src/airflow/sdk/io/path.py#L484-L500),这正是它能安全穿越 XCOM 序列化/反序列化流程的原因。

配置:连接机制与替代后端

连接配置自动下推

基本使用场景下,对象存储抽象几乎不需要额外配置,它完全依赖 Airflow 标准的Connection机制:通过conn_id指定要使用的连接,连接上的任何设置都会被下推到底层实现。例如使用 S3 时,可以在 Connection 中配置aws_access_key_idaws_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 Pathlibupath)之上,因此大部分操作与你操作本地文件系统的方式一致。本节只列出与标准 Path API 存在差异的操作,其余细节可查阅ObjectStoragePath类文档(对应实现见 task-sdk/src/airflow/sdk/io/path.py)。

mkdir

在指定路径或桶/容器内创建目录条目。对于没有真正目录概念的系统,可能仅为当前实例创建目录条目、并不影响真实文件系统。若parentsTrue,缺失的父级路径会被一并创建。

touch

在给定路径创建文件或更新时间戳。truncate默认为True,即会截断文件;若文件已存在,当exists_ok为真时操作成功(并把修改时间更新为当前时间),否则抛出FileExistsError

stat

返回一个类似stat_result的对象,支持st_sizest_mtimest_mode等属性,同时行为上又像一个字典,可提供对象的附加元数据。例如 S3 下会额外返回['ETag', 'ContentType']等键。如果代码需要跨对象存储移植,不要依赖这些扩展元数据。实现上,stat会把底层fs.stat结果包装成带protocolconn_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
containerbucket的别名path.py#L198-L200
fs便捷属性,返回已实例化的(Airflow 认证后的)文件系统path.py#L175-L178
key返回对象 key(按惯例去掉前导斜杠,保留尾部斜杠以支持目录语义)path.py#L208-L214
namespace返回对象的命名空间,通常是协议加桶名,如s3://bucketpath.py#L216-L218
path供文件系统实例使用的 fsspec 兼容路径继承自 UPath
protocolfsspec 协议名继承自 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,会确保读写落在紧随offsetoffset + length之后的分隔符边界上;若offset为 0 则从 0 开始;返回的字节串包含结尾分隔符;若offset + length超出 EOF,则一直读到 EOF。其 docstring 给出了一个直观例子(CSV 文件按换行符切块读取,见 path.py#L297-L311)。

复制与移动:跨对象存储的数据搬运

copymove用于把文件或目录从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;跨存储时退化为"先copyunlink"——这与前文"对象存储没有原子重命名"的局限一致。此外还提供了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 串成了真实的数据流:

  1. 模块顶层创建基础路径,conn_id直接内嵌在 URI 中:
base = ObjectStoragePath("s3://aws_default@airflow-tutorial-data/")
  1. 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 path
  1. analyze任务接收上游路径,注册path.fs到 DuckDB 并执行 SQL 查询(示例 L108-L129),完整演示了"写入对象存储 -> 通过 XCOM 传递路径 -> 外部引擎消费"的闭环。

深入实现:ObjectStoragePath 源码剖析

继承与认证文件系统注入

ObjectStoragePath继承自upath.extensions.ProxyUPath(path.py#L82),所有标准路径操作(existsmkdiriterdirglobwalkrenameread_byteswrite_bytes等)都委托给内部的__wrapped__UPath。构造时如果指定了 conn_id,实现会把 Airflow 认证后的文件系统直接注入__wrapped__._fs_cached,从而让所有委托操作"一次修复、处处生效",而不是为每个方法单独覆写;注入失败仅记录 DEBUG 日志,错误会在首次使用路径时才暴露(path.py#L148-L168)。

conn_id 的完整传播

conn_id会通过_from_upath从父实例传播到所有派生路径(/joinpathparentparentswith_namewith_suffixwith_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调用,自动向血缘收集器登记输入/输出资产;copymove在同存储或涉及本地文件时也会显式登记血缘。这意味着在使用对象存储 API 时,数据血缘追踪几乎是免费的。

向后兼容层

Airflow 核心包在 airflow-core/src/airflow/io/init.py 中通过add_deprecated_classesairflow.io.path.ObjectStoragePathairflow.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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 12:26:14

CANN/GE获取形状API文档

GetShape 【免费下载链接】ge GE(Graph Engine)是面向昇腾的图编译器和执行器,提供了计算图优化、多流并行、内存复用和模型下沉等技术手段,加速模型执行效率,减少模型内存占用。 GE 提供对 PyTorch、TensorFlow 前端的…

作者头像 李华
网站建设 2026/9/10 12:25:51

Python魔塔游戏源码解析:地图数据组织与战斗系统设计

简介:这是一份基于 Python 与 Pygame 开发的魔塔小游戏完整源码包,适合计算机相关专业学生用于毕业设计、课程设计或 Python 游戏开发入门练手。项目包含 main.py 主程序与 level.py 关卡配置文件,玩家可自由调整近乎所有关卡数据&#xff0c…

作者头像 李华
网站建设 2026/9/10 12:25:34

Magnitude不是CLI工具:轻量级向量相似度检索库解析

1. 项目概述:一个被严重误读的“magnitude”——它根本不是CLI工具,而是高性能向量相似度检索库 最近在多个技术社区和开发者群聊里,频繁看到有人把 magnitude 和各种 CLI 工具(codex cli、claude cli、grok cli)混为…

作者头像 李华
网站建设 2026/9/10 12:25:18

openGauss数据库安全架构与认证机制详解

1. openGauss数据库安全架构全景解析在企业级数据库领域,安全从来不是单一功能点的堆砌,而是贯穿整个系统生命周期的体系化工程。openGauss作为国产数据库的标杆产品,其安全设计采用了"纵深防御"理念,构建了四层立体防护…

作者头像 李华