Apache Airflow 修复 Backfill 提前标记完成问题:孤儿 Backfill 的基于时限的清理机制解析 Apache Airflow 修复 Backfill 提前标记完成问题孤儿 Backfill 的基于时限的清理机制解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读本文围绕 Apache Airflow 3.x 中一项针对 Backfill 生命周期管理的 bugfix见 airflow-core/newsfragments/62561.bugfix.rst展开修复了 DagRun 尚未创建完成时 Backfill 就被提前标记为 completed 的问题并新增了基于创建时限的孤儿 Backfill 自动清理机制。读完本文你将理解 Backfill 完成状态判定的完整链路、2 分钟初始化时间窗的设计动机、调度器定时扫描的实现细节以及如何通过单元测试验证该行为。问题背景Backfill 为何会被提前标记完成在 Apache Airflow 中Backfill 是为历史时间区间批量补跑 DAG 运行的核心机制。其状态流转由调度器Scheduler负责其中最关键的一步是当一个 Backfill 下所有的 DagRun 都进入终态SUCCESS / FAILED 等时将Backfill.completed_at字段写入时间戳从而标记该 Backfill 完成。在修复之前_mark_backfills_complete的判定条件存在一个竞态窗口Backfill 记录被创建后、对应的 DagRun 尚未落库之前调度器扫描时发现没有任何未完成状态的 DagRun 关联到该 Backfill就会立刻将 Backfill 标记为完成。这会导致用户提交的 Backfill 请求返回后数据库中的 Backfill 已显示 completed但实际没有任何 DagRun 被创建由于同一 DAG 只允许一个进行中的 Backfill见 airflow-core/src/airflow/models/backfill.py 中AlreadyRunningBackfill的约束被误标记完成的 Backfill 会阻塞后续合法请求用户无法通过 UI / API 观察到真正的执行进度。修复核心_mark_backfills_complete的守卫条件修复后的逻辑位于调度器作业运行器 airflow-core/src/airflow/jobs/scheduler_job_runner.py 的_mark_backfills_complete方法。其查询条件可拆解为三个部分unfinished_states (DagRunState.RUNNING, DagRunState.QUEUED) now timezone.utcnow() # 2 分钟初始化时间窗 initializing_cutoff now - timedelta(minutes2) query select(Backfill).where( Backfill.completed_at.is_(None), # 守卫条件Backfill 必须至少有一条关联记录 # 否则视为仍在初始化中参见 issue #61375 or_( exists(select(BackfillDagRun.id).where(BackfillDagRun.backfill_id Backfill.id)), Backfill.created_at initializing_cutoff, ), # 不存在任何处于 RUNNING / QUEUED 的 DagRun ~exists( select(DagRun.id).where( and_(DagRun.backfill_id Backfill.id, DagRun.state.in_(unfinished_states)) ) ), )三个判定条件的语义条件作用修复点completed_at IS NULL只处理尚未完成的 Backfill幂等性保障存在BackfillDagRun关联或created_at早于 2 分钟前确保 Backfill 已真正完成初始化或已超过初始化宽限期本次修复新增不存在 RUNNING / QUEUED 的 DagRun只有当所有 DagRun 均进入终态时才允许完成原有逻辑第一个or_分支是本次修复的关键只要 Backfill 下还没有任何BackfillDagRun关联记录就认为它仍处于初始化中源码注释明确指出这一守卫用于规避 issue #61375 描述的问题从而不会被提前标记完成。孤儿 Backfill 的基于时限清理机制单纯增加必须有 BackfillDagRun 关联的守卫会引入另一个问题如果初始化过程中途失败例如数据库锁冲突、进程崩溃Backfill 永远拿不到任何 DagRun 关联就会变成永久滞留的孤儿记录始终占用同一 DAG 只能有一个进行中 Backfill的配额。为此修复同时引入了基于创建时限age-based的清理机制initializing_cutoff now - timedelta(minutes2)即 Backfill 创建时间距今超过2 分钟且仍无任何关联记录时判定为初始化失败的孤儿允许被标记为完成completed_at now从而释放配额、允许用户重新发起 Backfill。for b in backfills: b.completed_at now这一设计的合理性在于正常情况下 Backfill 的初始化创建 Backfill 记录 批量写入 DagRun / BackfillDagRun在秒级内完成2 分钟的宽限期远大于正常耗时只有异常路径才会让 Backfill 在无任何关联的情况下存活超过 2 分钟。调度器 30 秒定时扫描_mark_backfills_complete由调度器主循环中的定时器驱动在 airflow-core/src/airflow/jobs/scheduler_job_runner.py 中注册timers.call_regular_interval( 30, self._mark_backfills_complete, )即每 30 秒执行一次完成状态扫描。结合 2 分钟初始化窗口可知最坏情况下一个初始化失败的孤儿 Backfill 会在创建后约 2 分 30 秒内被清理并标记完成不会无限滞留。测试验证覆盖完整生命周期边界本次修复配套的单元测试位于 airflow-core/tests/unit/jobs/test_scheduler_job.py以dag_maker构造序列化 DAG、以time_machine模拟时间流逝精确覆盖了各边界场景测试用例场景断言test_mark_backfills_completed正常 BackfillDagRun 全部置为 SUCCESScompleted_at被写入test_mark_backfills_complete_skips_initializing_backfillBackfill 创建后尚未建立任何 BackfillDagRun 关联未到 2 分钟completed_at保持None不被提前完成test_mark_backfills_complete_cleans_orphan_after_cutoff无关联的孤儿 Backfill 创建超过 2 分钟时间旅行 3 分钟completed_at被写入孤儿清理生效test_mark_backfills_complete_keeps_old_backfill_with_running_dagruns超过 2 分钟的旧 Backfill 但仍有 RUNNING DagRuncompleted_at保持Nonetest_mark_backfills_complete_young_backfill_with_finished_runs创建不足 2 分钟但所有 DagRun 已 SUCCESScompleted_at立即写入test_mark_backfills_complete_multiple_independent两个独立 Backfill一个完成一个运行中仅完成的被标记运行中的保持原状特别值得关注的是test_mark_backfills_complete_skips_initializing_backfill它精确复现了本次修复的 bug 场景——手动创建一个没有任何 BackfillDagRun 关联的 Backfill调用_mark_backfills_complete后断言completed_at is None随后补上 SUCCESS 状态的 DagRun 与关联记录再次调用后断言完成。这验证了有关联才允许完成的守卫确实生效。关联防护初始化失败的即时清理除了定时器的事后清理源码中还存在一条即时清理路径当并发创建 Backfill 触发数据库锁不可用错误is_lock_not_available_error时airflow-core/src/airflow/models/backfill.py 会调用_cleanup_partial_backfill做**尽力而为best-effort**的部分清理def _cleanup_partial_backfill(backfill: Backfill, session: Session) - None: Best-effort removal of a partially-created backfill after a lock error. try: session.rollback() session.execute(delete(BackfillDagRun).where(BackfillDagRun.backfill_id backfill.id)) session.execute(delete(DagRun).where(DagRun.backfill_id backfill.id)) session.delete(backfill) session.commit() except Exception: session.rollback()该函数删除残留的BackfillDagRun、DagRun与Backfill记录后提交若清理本身失败则回滚并让原始错误继续抛出调用方返回 503。它与 2 分钟时限清理互为补充即时清理处理创建即失败的异常时限清理兜底创建后无声失败的孤儿。运维启示不要手动修改completed_atBackfill 完成状态完全由调度器_mark_backfills_complete每 30 秒自动判定手工干预容易破坏幂等性与配额约束。关注初始化失败的信号若某个 DAG 频繁出现Backfill 提交后无任何 DagRun应检查数据库锁竞争与调度器日志而非反复重试同一 Backfill。了解 2 分钟宽限期的存在刚提交的 Backfill 在 UI 上短暂显示无运行是正常现象不要误判为失败2 分钟是孤儿判定的硬性边界。总结本次 bugfixairflow-core/newsfragments/62561.bugfix.rst通过必须有 BackfillDagRun 关联或超过 2 分钟创建时限的双重守卫同时解决了 Backfill 提前标记完成与孤儿 Backfill 滞留两个相互制约的问题前者保护了初始化窗口内的正确性后者保证了失败场景下的自愈能力。配合 30 秒定时扫描与完整的边界测试矩阵Backfill 生命周期管理在竞态与异常路径下均具备了可预期的行为。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考