Airflow Dag Serialization(DAG 序列化)完全指南:从无状态 Webserver 架构到 3.1+ 多语言 SDK 默认值契约
Dag Serialization 是 Apache Airflow 让 Webserver 摆脱 DAG 源文件依赖的关键机制:Scheduler 解析 DAG 后将其序列化为 JSON 写入 Metadata DB,Webserver 与调度决策一律从数据库读取序列化结果。本文基于 dag-serialization.rst 官方文档,并结合 serialized_dag.py、dagcode.py、renderedtifields.py 等源码实现,完整讲解其架构原理、三项核心配置项、底层读写链路,以及 Airflow 3.1 起在序列化 JSON 中引入的 client_defaults 版本化契约——读完你将掌握 DAG 序列化的运行机制、调优参数含义,以及如何让 Webserver 彻底无状态化。
一、为什么需要 DAG 序列化:让 Webserver 无状态、让调度更一致
在没有 DAG 序列化与 DB 持久化的时代(Airflow 1.10.x 早期),Webserver 和 Scheduler 都必须直接访问 DAG 文件,各自独立解析 DAG。这意味着每次 Webserver 启动、每次页面请求背后的 DagBag 构建,都要重新 import 全部 DAG 文件,导致启动慢、内存占用高,且 Webserver 与 DAG 文件系统强耦合。
Dag Serialization 的目标正是把 Webserver 从 DAG 解析中解耦出来,使其变得非常轻量:
- Scheduler 侧解析:由 Scheduler 内的
DagFileProcessorProcess(见 processor.py)解析 DAG 文件,把 DAG 序列化为 JSON 格式,保存到 Metadata DB 的SerializedDagModel(serialized_dag表,见 serialized_dag.py)中。 - Webserver 侧读取:Webserver 不再重复解析 DAG 文件,而是读取 DB 中的序列化 DAG JSON,反序列化后重建用于 UI 展示的 DagBag。
- 调度侧一致化:从 Airflow 2.0.0 起(配合 Scheduler HA),Scheduler 也不再依赖真实 DAG 文件做调度决策,而是使用包含全部调度所需信息的序列化 DAG,确保多 Scheduler 之间看到完全一致的 DAG 视图。
- 按需加载:Dag Serialization 的一项关键设计是——Webserver 启动时不再加载整个 DagBag,而是按需从
serialized_dag表逐条加载某个 DAG。当 DAG 数量庞大时,这一设计可显著减少 Webserver 的启动时间与内存占用。
注意:从 Airflow 2.0+ 起,Dag Serialization 是强制开启的,无法关闭。
从当前仓库源码可看到读取侧的实现细节:Webserver 侧使用 dagbag.py 中定义的 DBDagBag 从数据库取 DAG,它内部维护一个按 dag_version_id 索引的内存缓存 _CacheEntry(缓存反序列化后的 DAG 与其 dag_hash)。每次命中缓存时都会用 [core] min_serialized_dag_update_interval 作为"再验证节流窗口"(见 dagbag.py),在窗口期内直接服务缓存副本而不做 DB 往返,从而把最大数据陈旧程度限制在一个到两个更新周期内。
1.1 三个相互配合的数据模型
除了 SerializedDagModel,Dag Serialization 体系还涉及另外两张表:
DagCode(dag_code表):存储 DAG 源文件的完整代码,让 Webserver 完全独立于 DAG 文件系统。如果 DAG 文件已内嵌进 Docker 镜像,或者你能以其他方式提供给 Webserver,则不强制使用它。模型定义见 dagcode.py,其中source_code保存源码文本、source_code_hash保存 MD5 校验值,并通过dag_version_id与DagVersion关联;Scheduler 处理 DAG 时会调用write_code/update_source_code同步代码入库。RenderedTaskInstanceFields(渲染字段表):开启序列化后,模板字段不会在请求时即时渲染,而是在任务真正在 worker 上执行前,把字段内容的一份副本保存下来(模型见 renderedtifields.py)。为抑制数据库无限膨胀,只保留最近的若干条记录,更旧的记录会在任务执行期间被清理(详见下文配置项num_dag_runs_to_retain_rendered_fields)。
二、Dag Serialization 配置项详解(airflow.cfg)
在 airflow.cfg 的 [core] 段添加如下设置:
[core]
# You can also update the following default configurations based on your needs
min_serialized_dag_update_interval = 30
num_dag_runs_to_retain_rendered_fields = 30
compress_serialized_dags = False
三个配置项的作用如下:
| 配置项 | 默认值 | 作用说明 | 仓库证据 |
|---|---|---|---|
min_serialized_dag_update_interval |
30(秒) |
设定 DB 中序列化 DAG 允许被更新的最小时间间隔,用于降低数据库写入频率 | config.yml;collection.py 中 conf.getint("core", "min_serialized_dag_update_interval", fallback=30) |
num_dag_runs_to_retain_rendered_fields |
30(次运行) |
控制保留 Rendered Task Instance Fields 的最近 DAG 运行次数;更早运行的记录会在任务执行时被删除 | config.yml(默认 "30");renderedtifields.py |
compress_serialized_dags |
False |
是否在写入数据库前压缩序列化 DAG。集群中存在超大 DAG 时很有用;为 True 时会禁用 DAG dependencies(依赖)视图 |
config.yml;serialized_dag.py |
如果从早于 1.10.7 的版本升级,请务必先执行 airflow db migrate,以创建 serialized_dag、dag_code、rendered_task_instance_fields 等新表。
2.1 源码视角:压缩选项如何改变存储结构
compress_serialized_dags 的开关在底层直接影响表的存储列选择。在 serialized_dag.py 中,serialized_dag 表同时定义了两种列:
data:JSON 类型列(PostgreSQL 上使用JSONB变体);data_compressed:LargeBinary二进制列。
模型初始化时(见 serialized_dag.py):当 compress_serialized_dags=True 时,模块级常量 _COMPRESS_SERIALIZED_DAGS 为真,代码会对 sort_keys=True 排序后的 JSON 执行 zlib.compress 写入 data_compressed,同时把 data 置空;否则直接把字典存入 data 列。这也是该选项开启后 DAG 依赖视图被禁用的直接原因——依赖关系解析需要可读的 JSON 结构。
另外,SerializedDagModel 类文档字符串(serialized_dag.py)明确指出该表由 Scheduler 同步维护,除了上述两个 [core] 参数外,还受 [dag_processor] refresh_interval(默认 300 秒,用于清理已删除文件对应的旧序列化记录,官方建议在文件删除场景下调小到 60)影响。
2.2 写入与校验机制:dag_hash 与版本化
serialized_dag 表用 dag_hash(MD5)判断一个 DAG 是否真的发生了变化。在 serialized_dag.py 的 hash() 中,序列化数据会先被递归排序(保证相同内容产生相同哈希),再剔除 fileloc 与 bundle_name 后计算哈希——因为在 3.0+ 中,bundle_path + 相对 fileloc 的组合才真正决定 DAG 文件位置,单纯路径变化不应触发重写。
写入侧完整链路见 collection.py:DagFileProcessorProcess 完成解析后调用 _serialize_dag,其中读取 min_serialized_dag_update_interval 并传给 SerializedDagModel.write_dag(serialized_dag.py)——只有超过最小更新间隔且哈希变化时才会真正写库;若 DAG 未发生更新,则转而检查/更新 DagCode.update_source_code。读取侧则通过 SerializedDagModel.get_dag / get(serialized_dag.py)按 dag_id 取数。
三、渲染字段(Rendered Task Instance Fields)的保留与清理
启用序列化后,模板渲染从"请求时即时执行"变为"任务执行前在 worker 上保存一份字段内容副本"。这份副本存入 RenderedTaskInstanceFields 表,供 Webserver 的 Rendered View 标签页展示准确的渲染结果。
数据库膨胀控制是实现层面重点考虑的问题:delete_old_records()(renderedtifields.py)以 [core] num_dag_runs_to_retain_rendered_fields 为阈值,逻辑上先按 run_after 倒序取出该 DAG 最近 N 次运行的 run_id 列表(为避免全表扫描不直接扫 RTIF 表),再删除不在保留列表内的旧记录;被删除时以 dag_run 为粒度整体保留或整体删除某个运行内所有 mapped 子任务记录,避免出现残缺数据。MySQL 方言下由于不支持 LIMIT 子查询,会先取出 ID 列表再删除(见 renderedtifields.py),并且整个删除查询带有事务重试装饰器以应对可能的死锁。该清理时机在 taskinstance.py——任务实例写入渲染字段后即调用 delete_old_records(self.task_id, self.dag_id, ...)。
提示:把该数值调得过小可能导致查看较老任务的 Rendered 标签时报错,因为对应记录已被清理(见 config.yml)。
四、限制与规避:自定义过滤器 / 宏的渲染视图问题
当使用用户自定义过滤器(user-defined filters)和宏(macros)时,Webserver 的 Rendered View 对尚未执行的任务实例(TI)可能显示错误结果。原因在于:渲染依赖的外部模块位于 DAG 文件环境中,而无状态的 Webserver 未必能访问这些模块。
规避手段:
- 用
airflow tasks renderCLI 命令(命令定义见 cli_config.py)调试或测试template_fields的渲染结果; - 一旦任务真正开始执行,Rendered Template Fields 会按上文所述存入数据库的独立表(
RenderedTaskInstanceFields),此后 Webserver 的 Rendered View 标签页就会显示正确值。
提示:需要 Airflow >= 1.10.10 才能获得完全无状态的 Webserver;1.10.7 至 1.10.9 在部分场景下仍需要访问 DAG 文件。
五、替换默认 JSON 库(如 ujson)
序列化默认使用 Python 标准库 json——在 settings.py 中,settings.json 变量被显式声明为 "The JSON library to use for DAG Serialization and De-Serialization",而 serialized_dag.py 等模块正是通过 from airflow.settings import json 引用它(见 serialized_dag.py)。
要换成 ujson 这类更快的 JSON 实现,只需在本地 Airflow 设置文件 airflow_local_settings.py 中定义一个 json 变量并覆盖:
import ujson
json = ujson
Airflow 会通过 settings.import_local_settings()(见 settings.py)在启动早期导入 airflow_local_settings,从而让后续所有序列化/反序列化调用都走 ujson。关于本地设置文件的完整配置方式,参见 set-config.rst。
六、Airflow 3.1+:带默认值的序列化与 Task SDK 版本化契约
从 Airflow 3.1 开始,DAG 序列化确立了一份 Task SDK 与 Airflow 服务端组件(Scheduler、API-Server)之间的版本化契约。该契约与 Task Execution API 配合,将客户端与服务端组件解耦,使双方可以独立部署与升级,同时保持向后兼容并自动解析默认值。
6.1 默认值的生效顺序(优先级)
当 Airflow 处理 DAG 时,服务端按如下优先级依次应用默认值:
- Schema defaults(模式默认值):Airflow 内置默认值(优先级最低);
- Client defaults(客户端默认值):SDK 特定默认值;
- Dag
default_args(DAG 级设置):沿用既有行为; - Partial arguments(部分参数):MappedOperator 的共享值;
- Task values(任务显式值):显式设置的任务值(优先级最高)。
也就是说,你可以在不同层级设置默认值,越具体的设置会覆盖越通用的设置。
在源码侧,DagSerialization 类的 SERIALIZER_VERSION = 3(serialized_objects.py),反序列化时会读取 JSON 中的 __version 字段以选择兼容的解析路径。populate_operator、_apply_defaults_to_encoded_op、_matches_client_defaults 等逻辑(serialized_objects.py)共同实现"按层级补齐缺失值"的机制;generate_client_defaults()(serialized_objects.py)则在序列化端生成只包含"与 schema 默认值不同"条目的 client_defaults 段,并在单个 task 序列化时把与其匹配的字段省略掉,从而压缩 JSON 体积。
6.2 JSON 结构:新增 client_defaults 段
序列化后的 DAG 现在包含一个 client_defaults 段,用于存放常见的默认值:
{
"__version": 3,
"client_defaults": {
"tasks": {
"retry_delay": 300.0,
"owner": "data_team"
}
},
"dag": {
"dag_id": "example_dag",
"default_args": {
"retries": 3
},
"tasks": [{
"task_id": "example_task",
"task_type": "BashOperator",
"_task_module": "airflow.operators.bash",
"bash_command": "echo hello",
"owner": "specific_owner"
}]
}
}
6.3 值如何被应用
上例中,任务 example_task 的最终取值是:
- retry_delay:
300.0(来自client_defaults.tasks) - owner:
data_team(来自client_defaults.tasks,被任务内显式的owner: "specific_owner"?不——任务显式值拥有最高优先级,见下) - retries:
3(来自dag.default_args,覆盖client_defaults中可能存在的对应项) - bash_command:
"echo hello"(任务的显式值) - pool:
"default_pool"(来自 schema defaults)
系统会沿该层级自动为所有缺失字段补上值。需要特别说明:示例 JSON 中任务自身的 "owner": "specific_owner" 属于 Task values(第 5 级,最高优先级),因此该任务的最终 owner 实为 specific_owner,而非 data_team——这正是"更具体设置覆盖更通用设置"的直观体现。
6.4 MappedOperator(动态任务映射)的默认值处理
MappedOperators 同样参与默认值体系:
# Dag Definition
BashOperator.partial(task_id="mapped_task", retries=2, owner="team_lead").expand(
bash_command=["echo 1", "echo 2", "echo 3"]
)
在本例中,由 expand 生成的三个任务实例各自继承:
- retries:
2(来自 partial 参数) - owner:
"team_lead"(来自 partial 参数) - pool:
"default_pool"(来自 client_defaults,因为 partial 中未指定) - bash_command: 分别取
"echo 1"、"echo 2"、"echo 3"(来自 expand 的映射值)
6.5 独立部署架构与 SDK 合规要求
解耦后的组件划分:
- Server 组件(Scheduler、API-Server):负责编排,不运行用户代码;
- Client 组件(Task SDK、Dag processor):在隔离环境中运行用户代码。
关键收益:
- 独立升级:升级服务端组件无需触碰用户环境;
- 版本兼容:单一服务端版本可同时支持多个 SDK 版本;
- 部署灵活性:服务端与客户端组件可分别部署与扩缩容;
- 安全隔离:用户代码只运行在客户端环境,绝不在服务端组件上执行;
- 多语言 SDK 支持:任何语言都可以实现一个合规的 Task SDK(本仓库即包含 go-sdk、java-sdk、ts-sdk、task-sdk 等多语言实现)。
SDK 必须满足的要求:
- 遵循已发布的 schema:
- DAG 序列化:产出的 JSON 必须通过 dag-serialization schema 校验;
- 任务执行:通过 Execution API schema 支持运行期通信;
- 携带
client_defaults(可选):如有 SDK 特定默认值,放在client_defaults.tasks段中; - 使用正确的版本号:JSON 中必须带
__version字段标识序列化格式版本(当前为3)。
服务端承诺:
只要 SDK 同时遵守上述两份 schema 契约,Airflow 服务端组件就会:
- 正确反序列化来自任何合规 SDK 的 DAG;
- 在运行期支持任务执行通信;
- 按优先级层级正确应用默认值;
- 跨 SDK 版本与语言维持兼容性。
6.6 实现状态与演进方向
当前状态(Airflow 3.1): 序列化契约已为客户端/服务端解耦奠定基础。尽管部分服务端组件仍包含 Task SDK 代码(反之亦然),该契约已确保:
- 组件分离后,schema 合规性使独立部署成为可能;
- 无论代码耦合与否,版本兼容性始终成立;
- 部署分离在架构上已被支持,尽管尚未完全落地实现。
未来演进: 服务端与客户端组件之间的彻底代码解耦计划在未来版本完成;schema 契约提供的稳定接口,将在这一演进过程中持续保持一致。
七、总结与最佳实践建议
Dag Serialization 是贯穿 Airflow 调度与 Web 展示两个层面的基础设施:
- 生产环境应保持三个
[core]参数的默认基线(30/30/False),仅在遇到超大 DAG、写入压力大或数据库膨胀问题时针对性调整:写库频繁可调大min_serialized_dag_update_interval;超大 DAG 且不依赖 dependencies 视图时可开启compress_serialized_dags;渲染字段记录过多时按需缩小num_dag_runs_to_retain_rendered_fields(注意保留足够的调试回溯空间); - 若追求 Webserver 的彻底文件独立性,可依赖 Scheduler 自动同步的
dag_code源码存储,或将 DAG 文件打进镜像; - 升级自旧版本时勿忘
airflow db migrate; - 自定义宏/过滤器场景下,优先用
airflow tasks render验证渲染,再以任务真实执行后落库的渲染字段作为 UI 展示依据; - 在多语言 SDK / 独立部署演进中,始终以序列化 schema、
__version与client_defaults字段作为跨组件契约的核心。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0629
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python07
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
