Apache Airflow 渲染任务字段清理性能优化与 `num_dag_runs_to_retain_rendered_fields` 配置迁移指南
2026/9/10 8:18:25 网站建设 项目流程

Apache Airflow 渲染任务字段清理性能优化与num_dag_runs_to_retain_rendered_fields配置迁移指南

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

导读

在 Apache Airflow 中,DAG 序列化机制会把任务(Task)的模板字段在 Worker 上执行前预先渲染,并持久化到数据库的rendered_task_instance_fields(RTIF)表中,供 Web 界面 "Rendered" 标签页展示。当 DAG 中存在大量映射任务(Mapped Tasks)时,RTIF 记录会呈几何级数增长,历史记录的清理效率直接决定数据库负载。本文基于 airflow-core/newsfragments/60951.significant.rst 的发布说明,深入讲解本次针对映射任务场景的清理性能优化(约42 倍提速)、配置项max_num_rendered_ti_fields_per_tasknum_dag_runs_to_retain_rendered_fields的迁移,以及保留策略从"最近 N 次任务执行"到"最近 N 个 DAG Run"的语义变更,并结合源码与单元测试给出配置与验证实操。

背景:为什么 RTIF 清理性能至关重要

渲染字段的存储链路

启用 DAG 序列化(Airflow 2.0+ 强制开启)后,模板字段不再在 Web 请求时渲染,而是在任务执行前由 Worker 渲染一份副本存入数据库。其核心模型为RenderedTaskInstanceFields,表名为rendered_task_instance_fields,见 renderedtifields.py:

  • 主键由dag_idtask_idrun_idmap_index四列联合构成;
  • rendered_fields以 JSON 存储渲染后的模板字段;
  • k8s_pod_yaml可选存储 Kubernetes Executor 下渲染出的 Pod YAML;
  • 通过外键rtif_ti_fkeytask_instance表建立ON DELETE CASCADE关联。

每次任务实例写入渲染字段时,会调用 TaskInstance.update_rtif,其中rtif.write()负责原子 upsert(避免并发写冲突),随后立刻触发RenderedTaskInstanceFields.delete_old_records()进行历史记录清理。

映射任务带来的数据放大

映射任务是记录爆炸的放大器:一个映射任务的一次 DAG Run 会展开为多个map_index记录。单元测试 test_renderedtifields.py 明确验证了num_runs * 2条记录(两个映射实例)的写入场景。当 DAG 频繁运行且映射规模庞大时,RTIF 表会迅速膨胀,清理逻辑若对整表扫描或逐条删除,将拖慢每次任务执行。

核心优化:清理算法重写,性能提升约 42 倍

旧实现的问题

在本次变更之前,RTIF 清理基于"最近 N 次任务执行"进行保留,且清理过程需要扫描 RTIF 表本身来定位待删除记录。对于含大量映射任务的 DAG,每次任务执行都会触发一次昂贵的全表范围查询与逐条删除,产生显著的数据库开销。

新实现:以 DAG Run 为锚点

新的delete_old_records实现(见 renderedtifields.py)不再扫描 RTIF 表,而是直接从dag_run表按时间倒序取出最近 N 个run_id,再执行一条批量DELETE ... WHERE run_id NOT IN (...)

# 找到最近 N 个 dag run 的 run_id(不再扫描 RTIF 表) run_ids_to_keep_query = ( select(DagRun.run_id) .where(DagRun.dag_id == dag_id) .order_by(DagRun.run_after.desc()) .limit(num_to_keep) )

关键设计点:

  • run_after而非logical_date排序:源码注释明确指出logical_date在手动触发(manual runs)时可能为NULLrun_after作为实际的触发时间戳更可靠;
  • MySQL 兼容分支:由于 MySQL 不支持在IN/NOT IN子查询中使用LIMIT,代码会先物化出 run_id 列表(见 renderedtifields.py);
  • 批量删除_do_delete_old_records使用单条DELETE语句配合synchronize_session=False一次清除所有过期记录,并标注"此查询偶尔可能死锁,失败时由@retry_db_transaction装饰器自动重试"(见 renderedtifields.py)。

从测试断言可以印证性能提升的量化证据:test_delete_old_records参数化用例中,无论表中有 0、1、3、4、5 条 RTIF 记录,清理操作都只执行1 条查询expected_query_count=1,MySQL 因额外抓取 run_id 加 1 的 margin),见 test_renderedtifields.py。此前按执行次数保留并逐条删除的实现需要与记录数成正比的查询量,这正是"约 42 倍提速"的来源。

配置迁移:max_num_rendered_ti_fields_per_tasknum_dag_runs_to_retain_rendered_fields

新旧配置对照

本次变更将配置项改名,语义也从"每个任务保留的渲染字段条数上限"调整为"保留渲染字段的最近 DAG Run 数量":

项目旧配置新配置
配置名max_num_rendered_ti_fields_per_tasknum_dag_runs_to_retain_rendered_fields
所属 Section[core][core]
语义按任务执行次数保留最近 N 条按 DAG Run 数量保留最近 N 个 Run 的所有记录
默认值30(见 config.yml)
引入版本3.2.0(version_added: 3.2.0
旧名兼容仍可用,但会输出弃用警告

新的默认配置声明于 config.yml:

num_dag_runs_to_retain_rendered_fields: description: | Number of recent dag runs for which Rendered Task Instance Fields are retained. Records from older runs are deleted during task execution. Keeping this number small may cause an error when you try to view ``Rendered`` tab in TaskInstance view for older tasks. version_added: 3.2.0 type: integer example: ~ default: "30"

配置说明里同时给出运维警示:将该值调得过小,在 Web UI 查看较旧任务实例的 "Rendered" 标签页时报错,因为对应记录已被清理。

如何在 airflow.cfg 中配置

升级到新版本后,请在airflow.cfg[core]段使用新配置名:

[core] # 保留最近多少个 dag run 的 Rendered Task Instance Fields num_dag_runs_to_retain_rendered_fields = 30

若你仍使用旧配置名max_num_rendered_ti_fields_per_task,Airflow 会继续读取它但打印弃用警告(deprecation warning),建议尽快迁移。该配置同时可参考 dag-serialization.rst 中与其他序列化相关配置(min_serialized_dag_update_intervalcompress_serialized_dags)的组合使用示例。

行为语义变更:按 DAG Run 保留对稀疏任务的影响

从"执行次数"到"Run 数量"

旧逻辑保留的是最近 N 次任务执行的记录;新逻辑保留的是最近 N 个 DAG Run中该任务的全部记录。对于每个 Run 都执行一次的任务,两者结果接近;但对于条件任务/稀疏任务(conditional/sparse tasks)——即并非每个 Run 都会执行的任务——行为差异显著:

  • 新逻辑下,只要某个较旧的 Run 落出了"最近 N 个 Run"窗口,即使该任务的执行次数还不到 N 次,其记录也会被删除;
  • 结果就是"保留的记录数可能少于 N",这是官方发布说明中明示的预期行为("which may result in fewer records retained for conditional/sparse tasks")。

测试用例的精确验证

test_delete_old_records_sparse_task 完整演示了这一语义:

  1. 创建 10 个 DAG Run(run_0run_9),但仅在run_0run_3run_6run_9四个 Run 中写入 RTIF 记录(该任务是稀疏的,每 3 个 Run 才执行一次);
  2. 调用delete_old_records(num_to_keep=5)
  3. 由于保留窗口是"最近 5 个 Run"(run_5run_9),其中只有run_6run_9存在记录,最终恰好保留2 条,断言{r.run_id for r in result} == {"run_6", "run_9"}

这印证了:保留与否取决于 DAG Run 的新旧,而非任务执行的先后

映射任务记录的原子性保留

另一个重要保证是:同一个 DAG Run 内该任务的所有map_index记录会一起保留或一起删除,不会留下残缺数据。test_delete_old_records_mapped(见 test_renderedtifields.py)验证:5 个 Run、每个 Run 2 个映射实例共 10 条记录,设置num_to_keep=2后,保留 2 个最近 Run 的全部记录,最终记录数为remaining_rtifs * 2(整 Run 成对保留)。

参数说明与调优建议

delete_old_records的方法签名(见 renderedtifields.py)支持按任务粒度覆盖全局配置:

@classmethod @provide_session def delete_old_records( cls, task_id: str, dag_id: str, num_to_keep: int = conf.getint("core", "num_dag_runs_to_retain_rendered_fields", fallback=0), *, session: Session = NEW_SESSION, ) -> None:
  • 默认值读取自[core]段的num_dag_runs_to_retain_rendered_fieldsfallback=0表示配置缺失时回退为 0;
  • num_to_keep <= 0时直接返回(即不执行任何清理,见 renderedtifields.py),可用来在极端情况下关闭自动清理;
  • 清理发生在每次任务写入 RTIF 之后(update_rtif调用链),属于任务执行路径上的同步开销,因此新算法的低查询数直接降低了每次任务执行的数据库延迟。

调优建议(基于配置说明与代码逻辑,非性能基准):

  • 对"渲染字段查看需求频繁、Run 频率低"的 DAG,可适当调大该值(如 60~90),避免历史任务在 "Rendered" 页报错;
  • 对"Run 频率极高、映射规模大、存储压力大"的场景,可调小(如 5~10),但需接受旧任务无法查看渲染字段;
  • 若要彻底关闭自动清理,可设为0,但需自行规划 RTIF 表的手动维护策略,否则表会无限增长。

升级注意事项

  1. 配置迁移:将airflow.cfg[core] max_num_rendered_ti_fields_per_task替换为num_dag_runs_to_retain_rendered_fields,或暂时保留旧名接受弃用警告;
  2. 行为变化:升级后稀疏任务的 RTIF 保留数量可能少于旧版本,这是设计预期,无需惊慌;如需保持"按执行次数保留"的行为,Airflow 3.x 起不再支持,请按 Run 粒度重新评估保留需求;
  3. 数据库兼容:MySQL 环境下清理会多一次 run_id 物化查询(源码与测试均单独处理),属正常行为;
  4. 若升级自较早版本且涉及元数据库结构变更,请执行airflow db migrate(参见 dag-serialization.rst)。

小结

本次发布说明(60951.significant.rst)对应了 Airflow 对渲染字段存储子系统的一次实质性重构:通过将清理锚点从 RTIF 表切换到dag_run表、以单条批量DELETE取代逐条删除,使含大量映射任务的 DAG 的清理性能提升约 42 倍(测试断言单次清理恒定 1 条查询);同时以num_dag_runs_to_retain_rendered_fields取代旧配置,并将保留语义统一为"最近 N 个 DAG Run",为稀疏任务和映射任务提供了可预期、可原子化的保留行为。部署升级时,请同步完成配置迁移,并根据业务 Run 频率与渲染查看需求重新评估该参数的取值。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询