Airflow如何仅在backfill场景下执行任务前置数据库清理
问题背景
现有DAG任务依赖结构如下:
->task2 task1 ->task4 ->task3
所有任务的执行结果都会写入PostgreSQL、MongoDB存储:
- 日常按天调度运行时,对应处理日期的库表内无存量数据,任务运行无异常
- 执行历史日期backfill操作(例如task2代码更新后需要重算刷新历史结果)时,对应日期的存量数据会触发
duplicateError类报错,目前只能手动删除待backfill任务对应的日期数据后再执行重跑操作
需求为给每个任务编写对应执行日期的数据清理脚本,在实际业务任务执行前自动运行清空对应数据,同时出于性能考虑: - 清理逻辑仅在backfill场景下执行
- 日常按天调度DAG时不触发清理
核心待确认点: - Airflow中是否存在内置标识可判断当前任务是backfill触发还是日常调度触发
- 是否支持配置类似
pre_task_execute_when_backfill_callable的回调函数,在任务定义时指定backfill场景下的前置执行逻辑
具体实现方案
Airflow没有直接暴露名为is_backfill的专用标记,也没有内置命名完全匹配的backfill前置回调参数,但可以通过原生能力完全实现需求,不需要二次修改Airflow源码。
1. Backfill场景判断方法
方法1:按运行类型精准判断(Airflow 2.x+ 推荐)
Airflow 2.x及以上版本中,任务上下文里的dag_run对象自带run_type属性,直接标记了本次DAG运行的触发类型,取值共三类:
scheduled:调度器按配置的调度规则自动触发的日常运行manual:用户通过WebUI、API手动触发的单次运行backfill:通过airflow dags backfill命令触发的补数运行
判断逻辑非常简单,直接读取该字段即可100%识别命令触发的backfill场景:
def is_backfill(**context): dag_run = context.get("dag_run") return dag_run.run_type == "backfill" if dag_run else False
如果需要把手动选择历史日期触发的重跑场景也纳入清理范围,可以把判断条件调整为dag_run.run_type in ["backfill", "manual"],再叠加逻辑日期与当前时间的差值校验,避免误删当天手动触发的测试运行数据。
方法2:按时间差判断(兼容Airflow 1.x全版本)
日常按天调度的任务,实际触发时间永远晚于「逻辑日期 + 调度间隔」,不会出现延迟数天运行的情况;而backfill补历史数据时,实际触发时间会远晚于逻辑日期对应的时间窗口。可以通过时间差做判断,不需要依赖版本特性:
from datetime import datetime, timedelta def is_backfill(**context): # 按天调度场景下,时间差阈值设为2天即可覆盖正常调度延迟,可根据自身集群调度延迟情况调整 exec_date = context["data_interval_start"] current_time = datetime.now(tz=exec_date.tzinfo) return (current_time - exec_date) > timedelta(days=2)
2. Backfill前置清理逻辑的配置方式
不需要等待官方提供专用的回调参数,两种原生方式即可实现和pre_task_execute_when_backfill_callable完全一致的效果:
方式1:封装通用任务装饰器
把backfill判断、清理逻辑封装成通用装饰器,所有业务函数加上装饰器即可自动生效,不需要侵入原有业务代码:
from functools import wraps def run_clean_when_backfill(clean_func): def decorator(func): @wraps(func) def wrapper(**context): if is_backfill(**context): exec_date = context["data_interval_start"].strftime("%Y-%m-%d") clean_func(date=exec_date) return func(**context) return wrapper return decorator # 业务函数使用示例 @run_clean_when_backfill(clean_func=clean_task2_pg_and_mongo_data) def task2_business_logic(**context): # 原有task2的业务代码,不需要做任何修改 pass
方式2:自定义带前置清理逻辑的任务基类
如果不想给每个函数加装饰器,可以重写任务基类的pre_execute方法,自定义专用任务类,在任务初始化时传入对应清理函数即可:
from airflow.operators.python import PythonOperator class BackfillAutoCleanPythonOperator(PythonOperator): def __init__(self, backfill_clean_callable=None, *args, **kwargs): super().__init__(*args, **kwargs) self.backfill_clean_callable = backfill_clean_callable def pre_execute(self, context): super().pre_execute(context) if is_backfill(**context) and self.backfill_clean_callable: self.backfill_clean_callable(**context) # 任务定义示例 task2 = BackfillAutoCleanPythonOperator( task_id="task2", python_callable=task2_business_logic, backfill_clean_callable=clean_task2_related_data # 传入task2专属的PG、Mongo数据清理函数 )
这种方式下,每个任务只需要在定义时传入自己对应的数据清理逻辑,日常调度运行时pre_execute里的判断不通过,不会执行清理操作,完全满足性能要求;backfill场景下会先执行清理再跑业务逻辑,不会再出现主键重复的报错。
内容的提问来源于stack exchange,提问作者mmway

