
Apache Airflow Backfill 回填机制完全指南CLI、REST API 与 UI 实操及源码级原理解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowBackfill回填是 Apache Airflow 中为 DAG 的历史日期批量创建运行Dag Run的核心能力用于数据补录、失败重跑与逻辑修复等场景。本文以 Airflow 官方核心文档为骨架结合本仓库中 CLI 参数定义、Backfill 模型与 REST API 路由的源码实现系统讲解重处理行为、并发控制、反向运行、Dry Run、分区时间表回填的完整用法与底层工作原理。读完本文你将掌握通过 CLI、REST API 和 UI 三种方式精准控制历史数据重处理的完整实战方案。什么是 BackfillBackfill 是指为某个 DAG 的过去日期创建运行Dag Run的机制。Airflow 通过 CLI 和 REST API 提供该能力你只需要提供Dag、起始日期start date和结束日期end dateAirflow 就会按照该 DAG 的调度时间表schedule / timetable在给定日期范围内为每个应调度的逻辑日期logical date创建运行。核心要点Backfill 只对基于时间调度的 DAG 有意义。对于没有基于时间调度的时间表timetable的 DAG例如NullTimetable、OnceTimetable、ContinuousTimetable、由 Asset 触发的时间表等回填在概念上不成立Airflow 会直接拒绝。从源码看这一约束在 airflow-core/src/airflow/models/backfill.py 中被强制校验_do_dry_run与_create_backfill都会检查dag.timetable.periodic若为非周期性调度则抛出DagNonPeriodicScheduleException同时检查dag.allowed_run_types是否包含DagRunType.BACKFILL_JOB若不包含则抛出DagRunTypeNotAllowed。也就是说一个 DAG 是否允许回填既取决于时间表是否周期性也取决于其在元数据中声明的允许运行类型。数据重处理行为Reprocess Behavior控制回填时目标日期范围内可能已经存在历史运行例如上次回填留下的成功/失败记录。Airflow 提供了三种**重处理行为reprocessing behavior**选项决定遇到已存在运行时的处理策略选项行为none如果该逻辑日期已存在运行则无论其状态如何都不再创建新运行failed如果已存在运行且状态为 failed则为该日期创建新运行completed如果已存在运行且状态为 completed 或 failed则为该日期创建新运行一个重要约束如果该逻辑日期的最新一次运行仍在运行中running或处于排队queued状态则无论选择哪种重处理行为Airflow 都不会再创建新的运行。这一规则避免了对同一逻辑日期的并发重复执行。源码中的判定逻辑上述规则在 airflow-core/src/airflow/models/backfill.py 的_get_dag_run_no_create_reason函数中精确实现def _get_dag_run_no_create_reason(dr, reprocess_behavior: ReprocessBehavior) - str | None: non_create_reason None if dr.state not in (DagRunState.SUCCESS, DagRunState.FAILED): non_create_reason BackfillDagRunExceptionReason.IN_FLIGHT elif reprocess_behavior is ReprocessBehavior.NONE: non_create_reason BackfillDagRunExceptionReason.ALREADY_EXISTS elif reprocess_behavior is ReprocessBehavior.FAILED: if dr.state ! DagRunState.FAILED: non_create_reason BackfillDagRunExceptionReason.ALREADY_EXISTS return non_create_reason可以看到逻辑完全对应文档描述运行状态既不是 SUCCESS 也不是 FAILED即 running / queued / 其他中间态→ 判定为IN_FLIGHT进行中不创建新运行状态为终态时再按reprocess_behavior决定none一律不重跑failed只对 FAILED 状态重跑completed对 SUCCESS 与 FAILED 都重跑。ReprocessBehavior枚举定义于 airflow-core/src/airflow/models/backfill.py对应 CLI 的三个取值failed/completed/none。重处理时对已有运行的处理方式从源码看不同时间表类型的 DAG 在重处理时策略不同非分区 DAG当判定需要重跑时走_handle_clear_run逻辑backfill.py先dag.clear()清空已有运行的 Task Instance再将该运行更新为BACKFILL_JOB类型、triggered_byBACKFILL并可选地写入新的dag_run_conf然后重新排队执行。分区 DAG走_create_backfill_dag_run_partitionedbackfill.py不清除旧运行而是创建一条新运行与旧运行并存旧运行作为历史记录保留_get_latest_dag_run_row_query会按start_date取最新运行因此新创建的回填运行会成为后续生效的那一条。与 depends_on_past 的联动校验在 airflow-core/src/airflow/models/backfill.py 的_validate_backfill_params中还有一条易被忽略的规则如果 DAG 中存在depends_on_pastTrue的任务则反向回填reverseTrue被禁止抛出InvalidBackfillDirection重处理行为为none或未指定时被禁止抛出InvalidReprocessBehavior必须显式指定completed或failed。这是因为depends_on_past任务的执行依赖上一个调度周期的运行结果必须按时间正序且确保旧运行被重跑才能保证依赖链正确。并发控制max_active_runs回填支持通过max_active_runs参数限制本次回填内同时处于活跃状态的 Dag Run 数量。需要特别强调的是Backfill 的max_active_runs独立于 DAG 自身配置的max_active_runs生效——两者互不覆盖、互不合并回填的并发上限只约束本次回填创建的运行。参数定义与默认值在 airflow-core/src/airflow/cli/cli_config.py 中CLI 参数定义为ARG_MAX_ACTIVE_RUNS Arg( (--max-active-runs,), typepositive_int(allow_zeroFalse), helpMax active runs for this backfill., )即该参数必须是正整数不允许为 0。而在底层Backfill数据模型airflow-core/src/airflow/models/backfill.py中max_active_runs列的默认值为10max_active_runs: Mapped[int] mapped_column(Integer, default10, nullableFalse)也就是说即使你在 CLI 或 API 中不显式指定一次回填的默认并发上限也是 10。运行顺序--run-backwards回填默认按照时间从旧到新升序依次创建运行。你可以通过 CLI 的--run-backwards选项让回填反向执行即最新的日期最先运行。CLI 参数在 airflow-core/src/airflow/cli/cli_config.py 中定义其帮助文本还补充了一个关键限制If set, the backfill will run tasks from the most recent logical date first. Not supported if there are tasks thatdepend_on_past.即如果 DAG 中存在depend_on_pastTrue的任务则不支持反向回填前文已述源码会抛出InvalidBackfillDirection。反向执行的实现见 backfill.py 的_get_info_list先通过dag.iter_dagrun_infos_between(from_date, to_date)获取日期范围内的所有 DagRunInfo若reverseTrue则整体list(reversed(...))后再逐条创建。Dry Run回填预演Dry Run预演是 CLI 提供的一项能力在不真正创建任何 Dag Run 的前提下打印出回填将要考虑创建运行的日期清单。需要强调的是Dry Run 输出的清单只是将被考虑的日期这些日期最终是否真的会创建运行取决于你实际执行回填那一刻所选的重处理行为以及该范围内已有运行的状态。如果某个日期在预演时还没有运行但正式执行前恰好有了新的运行结果就会与预演输出不同。Dry Run 的实际输出在 airflow-core/src/airflow/cli/commands/backfill_command.py 中可以看到 Dry Run 的实现细节CLI 会先打印全部回填参数dag_id、from_date、to_date、max_active_runs、reverse、dag_run_conf、reprocess_behavior等然后调用_do_dry_run获取候选运行信息最后以表格形式输出每一行的logical_date、partition_key、partition_date。另外注意Dry Run 模式同样会执行全部参数校验_validate_backfill_params因此日期范围非法、目标 DAG 不存在或不允许回填等问题在预演阶段就会暴露。通过 CLI 创建回填完整示例CLI 是创建回填最直接的方式。以下命令创建一次对tutorialDAG 的回填覆盖 2015-06-01 至 2015-06-07 的日期范围airflow backfill create --dag-id tutorial \ --from-date 2015-06-01 \ --to-date 2015-06-07 \ --reprocess-behavior failed \ --max-active-runs 3 \ --run-backwards \ --dag-run-conf {my: param}命令参数逐一说明定义均见 airflow-core/src/airflow/cli/cli_config.py参数必填说明--dag-id是要回填的 DAG 标识--from-date是回填窗口起始时间包含对分区 DAG 自动解释为分区日期范围--to-date是回填窗口结束时间包含对分区 DAG 自动解释为分区日期范围--reprocess-behavior否取值none/completed/failed默认none控制已存在运行的重处理策略--max-active-runs否本次回填的最大并发运行数正整数默认 10--run-backwards否最新日期最先执行不支持含depend_on_past任务的 DAG--dag-run-conf否JSON 格式的dag_run.conf将注入到创建的每个 Dag Run--dry-run否仅打印将考虑创建的运行日期不真正创建--run-on-latest-version否实验性使用最新 bundle 版本执行任务而非创建原始运行时的版本关于--dag-run-confCLI 实现backfill_command.py会先用json.loads解析该字符串若 JSON 非法会直接抛出ValueError: Invalid JSON in --dag-run-conf。解析后的配置会经过 DAG 参数的校验_validate_backfill_params中会执行dag.params.deep_merge(dag_run_conf).validate()若配置与 DAG 声明的params冲突会抛出InvalidBackfillConf。另外CLI 的create_backfill命令同时承担两种模式backfill_command.py普通模式调用_create_backfill真正创建回填任务若该日期范围内没有任何可创建运行的日期例如窗口完全落在 DAG 首个调度边界之前会抛出NoBackfillRunsToCreateCLI 打印黄色警告并以退出码 1 结束。通过 UI 创建回填除了 CLI你也可以在 Airflow UI 中完成回填步骤如下进入某个 DAG 的Details详情页面点击Trigger触发按钮在弹出的窗口中选择Backfill回填运行类型填写表单Date range日期范围设置回填窗口的 From 与 To 逻辑时间logical datetimesReprocess behavior重处理行为在Missing Runs缺失运行、Missing and Errored Runs缺失与出错运行、All Runs全部运行三者中选择Max active runs最大活跃运行数限制本次回填的并发运行数Run backwards反向运行勾选后先执行最近的时间区间Advanced Config高级配置可选提供 JSON 格式的dag_run.conf如果该 DAG 当前处于暂停状态你可以在同一窗口中直接Unpause取消暂停它。UI 表单中的三个重处理行为选项与 CLI 的对应关系为Missing Runs对应none只补缺失运行、Missing and Errored Runs对应failed补缺失 失败重跑、All Runs对应completed成功与失败都重跑。分区 DAGPartitioned Dag的回填对于使用**基于分区的调度时间表partition-based timetable**的 DAG例如CronPartitionTimetable回填命令的--from-date/--to-date参数含义会自动切换Airflow 自动检测 DAG 是否为分区式partitioned日期范围被解释为分区日期范围partition-date range窗口内的每个分区对应创建一个 Dag Run。示例回填my_partitioned_dag从 2026-02-18 到 2026-02-20 的所有分区airflow backfill create --dag-id my_partitioned_dag \ --from-date 2026-02-18 \ --to-date 2026-02-20分区回填的并发同样由通用的--max-active-runs选项控制。分区与非分区的时间窗口差异在源码 backfill.py 的_get_info_list中两种模式对日期范围的处理存在一个细微差别非分区 DAG会额外过滤掉data_interval.end 当前时间的候选运行即不会为尚未结束的数据区间创建运行防止对未来数据回填分区 DAG直接使用iter_dagrun_infos_between返回的全部分区不做上述过滤。分区回填在创建 Dag Run 时使用时间表自身的generate_run_id并以partition_key、partition_date作为运行的身份标识_get_latest_dag_run_row_querybackfill.py在查找已有运行时会优先按partition_date分区刻度的 UTC 时刻匹配而非按partition_key字符串匹配这样即使分区键的格式化方式在时间表时区中发生变化重新命名也能正确识别同一分区已有的运行避免产生重复运行。REST API编程方式触发与查询回填除了 CLI 与 UIAirflow 还通过 REST API 暴露完整的回填能力路由实现在 airflow-core/src/airflow/api_fastapi/core_api/routes/public/backfills.py端点方法用途/public/backfillsGET分页列出指定 DAG 的回填记录/public/backfillsPOST创建回填/public/backfills/dry_runPOST回填预演Dry Run返回将创建的运行列表/public/backfills/{backfill_id}GET查询单个回填的详情/public/backfills/{backfill_id}/dag_runsGET列出该回填关联的 Dag Run包括被跳过的槽位/public/backfills/{backfill_id}/pausePUT暂停/恢复回填不会暂停已创建的运行REST API 的请求/响应数据模型定义在 airflow-core/src/airflow/api_fastapi/core_api/datamodels/backfills.py其中reprocess_behavior与 CLI 共用同一个ReprocessBehavior枚举。对应的 API 行为测试见 airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_backfills.py覆盖了创建回填、Dry Run、非法配置拒绝、数据库锁竞争、空窗口返回空列表等场景可作为理解接口语义的参考。回填的底层数据模型与状态机最后从数据模型层面理解一次回填的完整生命周期。核心表结构定义在 airflow-core/src/airflow/models/backfill.pybackfill表回填任务本体包含dag_id、from_date、to_date、dag_run_confJSON 列、reprocess_behavior默认none、max_active_runs默认 10、is_paused、triggering_user_name、created_at、completed_at、updated_atbackfill_dag_run表回填与 Dag Run 的关联表记录每个逻辑日期/分区的处理结果若某日期因故未创建运行dag_run_id为空并写入exception_reason可取值为in flight运行中/排队中、already exists按重处理行为无需重跑、unknown。一次回填创建后Airflow 会遍历日期范围内的每个调度点逐个执行查询最新运行 → 按重处理行为判定 → 创建新运行 / 清空重跑 / 跳过的流程并通过行级锁with_row_locks(..., skip_lockedTrue)防止并发场景下的重复创建。如果同一 DAG 已存在未完成completed_at IS NULL的回填新的回填请求还可能被拒绝AlreadyRunningBackfill从而保证同一 DAG 不会同时有多个活跃回填任务互相干扰。总结Backfill 是 Airflow 历史数据治理的核心工具其设计可以用三个维度 三个入口来概括三个维度重处理行为none/failed/completed决定哪些已有运行要重跑max_active_runs决定一次回填多快--run-backwards决定从新到旧还是从旧到新三个入口CLIairflow backfill create、REST API/public/backfills系列端点、UIDag Details 页的 Trigger → Backfill 弹窗三者共享同一套校验逻辑与ReprocessBehavior枚举行为完全一致边界条件非周期性调度的 DAG 不支持回填运行中/排队中的最新运行不会被重跑depend_on_past任务禁止反向回填且必须显式指定重处理行为分区 DAG 的日期范围自动按分区解释。掌握了这些参数语义与源码层面的判定规则你就能在数据补录、失败重跑、代码逻辑变更后的历史数据重算等场景中精确、可控地执行回填而不会产生重复或遗漏。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考