Apache Airflow 3 DAG Bundles 实战指南:用本地目录、Git、S3 与 GCS 实现 DAG 的版本化交付
Apache Airflow 3 DAG Bundles 实战指南用本地目录、Git、S3 与 GCS 实现 DAG 的版本化交付【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本指南围绕 Apache Airflow 3 引入的Dag BundleDAG 包机制展开系统讲解它相较于 Airflow 2dags_folder的核心演进、四种内置 Bundle 类型Local / Git / S3 / GCS的配置方法、重跑版本行为rerun_with_latest_version的控制逻辑以及如何基于BaseDagBundle编写自定义 Bundle。读完本文你将能在一套 Airflow 部署中同时挂载多种 DAG 来源并为 DAG Run 实现同一份代码跑完整次运行的可复现交付。DAG Bundle 是什么为什么它是 Airflow 3 的基石Dag Bundle 是一个或多个 DAG 文件及其关联资源辅助 Python 脚本、配置文件等的集合它可以来自本地目录、Git 仓库或 S3/GCS 等外部系统。管理员可以为一套 Airflow 部署定义多个 Bundle让 DAG 按逻辑单元组织管理由于 Bundle 处于包这一更高抽象层级天然支持对 DAG 运行所需的一切内容做整体版本化。Airflow 2 及更早版本中DAG 必须集中放在本机磁盘的单一dags_folder里把代码搬进去完全依赖部署管理员手工操作。DAG Bundle 与这种Dags folder模式相似但更强大版本控制Version ControlBundle 支持版本化后DAG Run 在整个运行期间都使用同一份代码——即使运行中途 DAG 被更新本次运行也不受影响重跑时仍可选择原始版本复现问题可扩展性Scalability把大量 DAG 组织成逻辑单元Airflow 能更高效地管理它们灵活性Flexibility无缝对接 Git、对象存储等外部系统来获取 DAG 源码摆脱必须放到本机固定目录的束缚。从源码结构看整个机制的核心抽象是 BaseDagBundle 基类所有 Bundle 类型无论内置还是自定义都实现这套接口DAG Bundle 同时被DAG Processor用于持续解析并保持 DAG 最新与Worker用于按某个特定版本运行任务两种角色使用二者使用模式不同Processor 始终只跟踪最新版本而单个 Worker 上可能同时存在同一 Bundle 的多个版本副本。内置 Bundle 类型一览Bundle 类型classpath版本化支持适用场景本地目录airflow.dag_processing.bundles.local.LocalDagBundle否始终用磁盘上最新代码开发、测试环境向后兼容 Airflow 2 的dags_folderGit 仓库airflow.providers.git.bundles.git.GitDagBundle是记录 Git commit生产环境 GitOps 化交付 DAGS3airflow.providers.amazon.aws.bundles.s3.S3DagBundle否始终用桶内最新代码以 S3 对象存储分发 DAGGCSairflow.providers.google.cloud.bundles.gcs.GCSDagBundle否始终用桶内最新代码以 GCS 对象存储分发 DAG版本化能力通过基类上的类属性supports_versioning: bool False声明LocalDagBundle在 local.py 中显式置为FalseGitDagBundle在 git.py 中置为True。配置 DAG Bundlesdag_bundle_config_list所有 Bundle 都在配置项[dag_processor] dag_bundle_config_list中声明可以在此配置一个或多个 Bundle。该配置项是一个 JSON 列表其中每个元素是一个字典包含name、classpath、kwargs三个键默认值为一个指向 Dags 文件夹的LocalDagBundle以保持与 Airflow 2 一致的开箱行为[dag_processor] dag_bundle_config_list [ { name: dags-folder, classpath: airflow.dag_processing.bundles.local.LocalDagBundle, kwargs: { path: /opt/airflow/dags } } ]对LocalDagBundle而言唯一的 kwarg 是path省略时默认取[core] dags_folder的值。这一定义与 config.yml 中记录的默认配置完全一致default: - [ {{ name: dags-folder, classpath: airflow.dag_processing.bundles.local.LocalDagBundle, kwargs: {{}} }} ]配置加载与合法性校验源码级在 DagBundlesManager.parse_config 中配置经过严格的解析与校验配置值必须是一个列表否则抛出AirflowConfigException每个元素必须是字典并通过 Pydantic 模型_ExternalBundleConfig含name、classpath、kwargs、可选的team_name校验Bundle 名称不允许重复重复会直接报配置错误名称example_dags是保留名称不可占用示例 DAG 应通过[core] load_examples配置启用当[core] load_examples为True时管理器会自动追加example_dagsBundle并把每个已安装 Provider 的example_dags目录注册为名为provider-example-dags的本地 Bundle见 _add_provider_example_dags_to_bundleclasspath 通过import_string动态导入类随后缓存在内存字典中同一类不会被重复导入启动时通过sync_bundles_to_db把配置同步为元数据库中的DagBundleModel记录从配置中移除的 Bundle 会被标记为active False对多 dag-processor 分片部署可传deactivate_missingFalse关闭该行为。凭据必须走 Connection严禁内联Bundle 的kwargs原样存储于[dag_processor] dag_bundle_config_list配置中而 Airflow 在启用[api] expose_config时会通过 Config API 对外暴露配置。任何有配置读取权限的用户都能逐字读取这些值因此kwargs中绝不能包含密钥。不要写出这样的配置repo_url: https://x-access-token:tokengithub.com/org/repo.git # 严禁正确做法是引用 Airflow 的 ConnectionGit 用git_conn_id、S3 用aws_conn_id、GCS 用gcp_conn_id把真实凭据放到 secrets backend 中Connection 的字段在运行时才解析不会落盘写入dag_bundle_config_list。Airflow 的 secrets backend 能力覆盖本地文件、环境变量与外部服务等多种实现具体可参见 secrets backend 文档。GitDagBundle版本化 DAG 的 GitOps 方案Git Bundle 是唯一内置支持版本化的对象存储/Git 类 Bundle运行每个 DAG Run 时会记录其创建时对应的 Git commit即使仓库后续已更新重跑仍使用完全相同的代码。kwargs 说明kwarg必填说明tracking_ref是分支名、tag 或 commit SHABundle 跟踪的目标引用git_conn_id否持有仓库凭据SSH/token的 Airflow Connection未传时回退为git_defaultrepo_url否显式指定仓库 URL可覆盖 Connection 的 host若传了 Connection 则二者二选一或配合使用subdir否仓库内存放 DAG 的子目录Bundle 的path会指向repo_path / subdirsparse_dirs否启用稀疏检出仅克隆指定目录列表需要 git ≥ 2.25refresh_interval否覆盖全局刷新间隔对应参数签名可直接在 GitDagBundle.init中确认。一个完整的配置示例[dag_processor] dag_bundle_config_list [ { name: my-git-repo, classpath: airflow.providers.git.bundles.git.GitDagBundle, kwargs: { git_conn_id: my_git_conn, subdir: dags, tracking_ref: main, } } ]底层实现要点git.pyrefresh()会根据tracking_ref拉取远端origin/tracking_refget_current_version()返回当前 checkout 的 commit SHA从而把代码版本与运行版本一一对应起来。Git Bundle 的完整 kwargs 列表与更多示例可参考apache-airflow-providers-git的 bundles 文档本仓库中对应实现位于 providers/git。S3DagBundle 与 GCSDagBundle从对象存储加载 DAG两者均以bucket_name为唯一必填 kwarg支持用prefix把 Bundle 限定到桶内子目录且均不支持版本化任务始终跑最新代码。S3 配置示例[dag_processor] dag_bundle_config_list [ { name: my-s3-dags, classpath: airflow.providers.amazon.aws.bundles.s3.S3DagBundle, kwargs: { aws_conn_id: aws_default, bucket_name: my-airflow-bucket, prefix: dags/ } } ]aws_conn_id默认值为aws_default。从源码可见s3.py初始化时会校验桶与 prefix 是否存在。完整 kwargs 参考apache-airflow-providers-amazon的 bundles 文档仓库实现位于 providers/amazon/src/airflow/providers/amazon/aws/bundles/s3.py。GCS 配置示例[dag_processor] dag_bundle_config_list [ { name: my-gcs-dags, classpath: airflow.providers.google.cloud.bundles.gcs.GCSDagBundle, kwargs: { gcp_conn_id: google_cloud_default, bucket_name: my-airflow-bucket, prefix: dags/ } } ]gcp_conn_id默认值为google_cloud_default。对应实现位于 providers/google/src/airflow/providers/google/cloud/bundles/gcs.py。在同一部署中组合多种 Bundle不同类型的 Bundle 可以在一个dag_bundle_config_list里自由混用默认的dags-folder本地 Bundle 既可删除也可保留与其他 Bundle 并存[dag_processor] dag_bundle_config_list [ { name: my_git_bundle, classpath: airflow.providers.git.bundles.git.GitDagBundle, kwargs: {tracking_ref: main, git_conn_id: my_git_conn} }, { name: dags-folder, classpath: airflow.dag_processing.bundles.local.LocalDagBundle, kwargs: {} } ]注意多行值的缩进尤其是最后一行必须保持正确否则 INI 多行值解析会失败——这是 Pythonconfigparser的既有行为详见其支持的 INI 文件结构文档。此外配置解析阶段还要求Bundle 名称在整个列表中唯一且不能使用保留名example_dags。为每个 Bundle 单独定制刷新频率全局刷新由[dag_processor] refresh_interval控制默认 300 秒决定 DAG Processor 多久去 Bundle 源里找新文件。可以通过 kwargs 传入refresh_interval对单个 Bundle 覆盖该值{ name: my-git-repo, classpath: airflow.providers.git.bundles.git.GitDagBundle, kwargs: { subdir: dags, tracking_ref: main, refresh_interval: 0 } }在源码层面BaseDagBundle.__init__的refresh_interval参数默认值即取自conf.getint(dag_processor, refresh_interval)见 base.pyconfig.yml 的示例配置也演示了按 Bundle 设 0表示尽快刷新的用法config.yml。自定义视图 URLview_url_template默认情况下每个 Bundle 类型会提供自己的视图 URL例如 Git Bundle 指向仓库 Web 界面。如需改用自定义 URL可在该 Bundle 的kwargs中传view_url_template[dag_processor] dag_bundle_config_list [ { name: my_git_repo, classpath: airflow.providers.git.bundles.git.GitDagBundle, kwargs: { tracking_ref: main, git_conn_id: my_git_conn, view_url_template: https://my.custom.git.repo/view/{subdir}, } } ]规则如下模板中的占位符如{subdir}只能是Bundle 自身拥有的属性渲染时会替换为对应属性值manager.py 的_extract_template_params通过正则\{([^}])\}提取占位符并取getattr(bundle, placeholder)不允许使用 Bundle 属性之外的任何占位符指定自定义 URL 后它会覆盖该 Bundle 默认提供的 URLURL 会经过安全检查_is_safe_bundle_urlmanager.py只允许http/httpsscheme、必须有网络位置netloc、不得包含控制字符若判定不安全该 Bundle 的视图 URL 会被置为None以防潜在安全问题。值得一提的是安全 URL 还会用[core] fernet_key做签名后落库signed_url_template用于完整性校验视图 URL 的生成不需要初始化 Bundle因此 UI 在展示时无需触发拉取等昂贵操作。旧版view_url方法已标记弃用请使用view_url_template。存储布局与磁盘清理理解底层运行方式尽管配置是逻辑层面的但理解 Bundle 在磁盘上的落点有助于排查问题相关辅助函数见 base.py存储根目录由[dag_processor] dag_bundle_storage_path指定必须为绝对路径未设置时默认落到Path(tempfile.gettempdir()) / airflow / dag_bundles每个 Bundle 的基础目录为root/bundle_name各版本副本存放在root/bundle_name/versions/version共享磁盘的 Worker 上Bundle 版本副本会随任务运行不断累积因此 Airflow 提供了一组清理配置见 config.yml配置项默认值含义stale_bundle_cleanup_interval1800两次陈旧 Bundle 清理检查的间隔秒数设为0或负数可禁用stale_bundle_cleanup_age_threshold21600距上次使用超过该秒数的版本才可能被删除stale_bundle_cleanup_min_versions10本地至少保留的版本数清理由BundleUsageTrackingManagerbase.py执行它跟踪每个版本最后被使用的时间戳删除使用不活跃且超过阈值的版本但总是保留最近使用的 N 个版本删除前通过文件锁flock确保不会误删正在被任务使用的副本——这也呼应了Worker 上多个版本并存、每个版本按需锁定使用的设计。Docker 镜像中安装 Git从Airflow 3.0.2起官方基础镜像已预装 git若你使用更早的版本或自建镜像需要自行安装 git 并配置gitpython使用它RUN apt-get update apt-get install -y git ENV GIT_PYTHON_GIT_EXECUTABLE/usr/bin/git ENV GIT_PYTHON_REFRESHquiet其中GIT_PYTHON_GIT_EXECUTABLE告诉 GitPython 去哪里找 git 可执行文件GIT_PYTHON_REFRESHquiet抑制 GitPython 每次调用的仓库刷新日志。在用户模拟run_as_user下使用 DAG Bundle当结合run_as_user用户模拟使用 Bundle 时需要保证被模拟的用户能访问由 Airflow 主进程创建的 Bundle 文件请按以下两点配置权限所有被模拟用户与 Airflow 运行用户属于同一个用户组配置合适的 umask例如umask 0002使组内成员可读/写新建文件。注意这种基于共享组权限的做法只是临时方案。未来的 Airflow 版本将通过 supervisor 化的 Bundle 操作来管理多用户访问届时无需再依赖共享组权限。控制重跑使用哪个版本rerun_with_latest_version当一个 DAG Run 或 Task Instance 被 clear 时UI 会弹出一个复选框询问重跑时使用最新 Bundle 版本还是原运行所用的版本。rerun_with_latest_version设置即控制该复选框的默认勾选状态团队无需每次手动决策该设置同样决定通过 API 或 CLI 创建 backfill 时的默认run_on_latest_version行为。工作原理每个 DAG 都有一个解析版本DagModel.bundle_version每次 DAG Processor 重新解析 DAG 文件时更新每个 DAG Run 则记录它创建时所使用的 Bundle 版本当rerun_with_latest_version为False时clear 一个 DAG Run 会保留其原始 Bundle 版本重跑使用同一份代码——这在排查失败时可复现为True时clear 会把 DAG Run 更新到当前解析版本确保重跑采用最新代码。解析优先级从高到低显式请求API 请求体中的run_on_latest_version参数若提供DAG 级DAG 构造参数rerun_with_latest_version显式传True/False时全局配置[core] rerun_with_latest_version选项若设置调用点回退clear/rerun 场景默认Falsebackfill 场景默认True保留各自的历史默认行为。一个例外自身没有任何版本的 DAG Run从 Airflow 2 迁移而来或版本被airflow db clean清理掉没有可保留的版本因此 clear 它时无论设置如何总是使用最新代码与最新 Bundle 版本。全局配置[core] rerun_with_latest_version False # 重跑时使用原始 Bundle 版本 # rerun_with_latest_version True # 重跑时使用最新 Bundle 版本未设置时按第 4 条调用点回退规则处理。该配置项在 config.yml 中定义version_added: 3.2.0。DAG 级配置可按 DAG 覆盖全局默认值from datetime import datetime from airflow import DAG from airflow.operators.empty import EmptyOperator # 该 DAG 总是用最新版本重跑 with DAG( dag_idalways_latest_dag, rerun_with_latest_versionTrue, start_datedatetime(2024, 1, 1), ) as dag: EmptyOperator(task_idtask)典型使用场景调试失败运行用False默认clear 失败运行后以同一份代码重跑便于复现与定位问题总是运行最新代码团队更希望重跑总是拿到最新代码时例如原运行后又部署了修复全局设[core] rerun_with_latest_version True混合策略全局默认设为True但对关键 DAG 用rerun_with_latest_versionFalse覆盖在最需要版本稳定性的地方保持固定版本。注意以上行为仅适用于支持版本化的 Bundle 类型如GitDagBundle。本地 Bundle 不支持版本化始终使用最新代码不受该设置影响。与 disable_bundle_versioning 的关系Airflow 提供两个影响 Bundle 版本行为的独立设置二者目的不同disable_bundle_versioning彻底关闭版本跟踪。置为True后DAG Run 上不再记录bundle_version。可同时作为 DAG 参数与全局配置[dag_processor] disable_bundle_versioning默认False见 config.yml使用且只对支持版本化的 Bundle 生效rerun_with_latest_version在保持版本跟踪开启的前提下控制默认的重跑行为仅改变呈现给用户的默认选择版本历史仍被完整记录。一句话总结disable_bundle_versioning回答是否要跟踪版本rerun_with_latest_version回答重跑时默认用哪个版本。两者相互独立且当版本化被禁用时rerun_with_latest_version不产生任何效果。编写自定义 DAG Bundle当内置类型无法满足需求如自建代码托管平台、企业内部对象存储时可继承BaseDagBundle实现自定义 Bundle。必须实现的抽象方法自定义类必须实现以下三个成员基类中用abstractmethod声明见 base.pypath属性返回一个Path指向本 Bundle DAG 文件存放的目录。Airflow 依据它定位并解析 DAG 文件initialize()之后该路径下应能访问到全部 DAG 文件get_current_version()返回当前 Bundle 版本的字符串。之后 Airflow 运行任务时会把该版本传回__init__以重新获取同一版本若不支持版本化则返回None。推荐返回BundleVersion数据类同时携带版本字符串与可选的结构化元数据如 S3 manifest直接返回裸字符串属于已弃用的旧路径未来会发出弃用警告base.pyrefresh()负责从源头刷新 Bundle 内容例如从远端拉取最新提交DAG Processor 会周期性调用它以保证 Bundle 不过期。可选覆写的方法__init__可扩展以接收额外参数如 Git Bundle 的tracking_ref必须调用父类__init__保证name、refresh_interval、version、version_data、view_url_template等基础字段正确初始化。避免在本方法中做网络调用等昂贵操作会拖慢 Bundle 实例化放到initialize中做initialize()在 DAG Processor 或 Worker首次真正使用Bundle 内容前被调用适合执行昂贵的准备工作它只在 Airflow 需要把 Bundle 文件落到磁盘时触发纯调用view_url的场景不会触发。覆写时务必在方法末尾调用super().initialize()基类实现会校验 Bundle 路径存在并给出告警base.pyview_url_template()返回一个字符串模板用于在 UI 中跳转到外部系统如 Git Web 界面查看该 Bundle 某版本它需要在未调用initialize的情况下也能工作。一个最小化的本地目录 自定义源风格示例from pathlib import Path from airflow.dag_processing.bundles.base import BaseDagBundle class MyStaticBundle(BaseDagBundle): supports_versioning False # 或不声明基类默认即为 False def __init__(self, *, source_path: str, **kwargs) - None: super().__init__(**kwargs) self.source_path Path(source_path) property def path(self) - Path: return self.source_path def get_current_version(self) - None: # 不支持版本化 → 返回 None return None def refresh(self) - None: # 例如从自定义源同步文件到 self.path ... def view_url_template(self) - str | None: return https://example.com/browse/{version}随后即可在dag_bundle_config_list中通过classpath指向它。其他关键注意事项版本化 Bundle若你的 Bundle 支持版本化请确保initialize、get_current_version、refresh都正确处理版本相关的逻辑按版本检出/保留副本并发安全Worker 可能同时创建大量 Bundle 对象Airflow不会序列化对 Bundle 对象的调用。如果底层技术对此敏感例如多个进程同时 clone 同一 git 仓库Bundle 类必须自行加锁保证同一时刻只有一个对象在克隆。基类提供了可复用的lock上下文管理器基于fcntl.flock的排他文件锁见 base.pyTriggerer 限制DAG Bundle 不会在 triggerer 组件中被初始化因此 trigger 代码不能来自 DAG Bundle——triggerer 不处理随时间变化的 trigger 代码一切都在主进程中发生。若需要自定义 trigger请确保它位于 Python 环境sys.path中而不是从 Bundle 获取。升级与迁移要点DAG Bundle 是 Airflow 3.0 的核心新能力元数据库中的dag_bundle表由迁移脚本引入见 0050_3_0_0_add_dagbundlemodel.py后续版本又加入了 URL/模板参数、bundle_name非空约束等演进默认dags-folderBundle 的存在保证了从 Airflow 2 平滑升级旧 DAG 仍然从dags_folder加载。若部署使用了自定义 Bundle历史遗留 DAG 行bundle_name指向未配置 Bundle 或relative_fileloc为空会在解析生命周期中自愈必要时可执行airflow dags reserialize强制重写bundle_name与relative_fileloc从 Airflow 2 迁移过来、本身不带 Bundle 版本的 DAG Runclear 时总是采用最新版本无旧版本可保留。通过将 DAG 交付升级为版本化的 BundleAirflow 3 让大规模团队协作下的 DAG 生命周期管理走向了 GitOps 化代码进仓库、运行按版本、重跑可复现、多源可并存。本文涉及的配置、类与源码入口都可在本仓库对应路径中进一步展开研读。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考