Airflow 2.0中如何用TriggerDagRunOperator触发目标DAG执行回填?
Airflow 2.x 触发DAG回填的配置方案
默认TriggerDagRunOperator只会触发与当前任务执行日期相同的单条DAG Run,要实现批量回填不受历史记录影响,可以用以下两种方案:
方案1:循环生成多Trigger任务(推荐,符合Airflow设计规范)
遍历回填时间范围内的所有日期,为每个日期生成独立的触发任务,显式指定执行日期,每个日期对应独立的DAG Run,状态独立可单独重试。
触发侧代码示例:
from airflow.operators.trigger_dagrun import TriggerDagRunOperator from datetime import datetime, timedelta # 替换为实际需要的回填起止日期 backfill_start_date = datetime(2021, 11, 15) backfill_end_date = datetime(2023, 10, 1) current_date = backfill_start_date trigger_tasks = [] while current_date <= backfill_end_date: trigger_task = TriggerDagRunOperator( task_id=f'trigger_target_{current_date.strftime("%Y%m%d")}', trigger_dag_id='target_dag', conf={'message': f'Running target for {current_date.strftime("%Y-%m-%d")}'}, execution_date=current_date, # 显式指定执行日期 reset_dag_run=True, # 覆盖同日期已存在的DAG Run,忽略历史运行记录 wait_for_completion=False, # 不需要等待单条完成可设为False,提高并发效率 ) trigger_tasks.append(trigger_task) current_date += timedelta(days=1) # 按需将trigger_tasks挂载到你的上游节点即可 # 示例:上游任务 >> trigger_tasks
该方案无需修改原target_dag的代码,且原target_dag配置的max_active_runs=10会自动限制并发数,最多同时跑10个日期的任务,不会超出资源限制。
方案2:单触发任务批量处理(适合小数据量短周期回填)
在触发时传入回填时间范围,修改目标DAG的任务逻辑,循环处理所有日期的数据。
触发侧代码修改:
trigger1 = TriggerDagRunOperator( task_id = 'trigger1', trigger_dag_id = 'target_dag', conf = { 'message': 'Starting target backfill', 'backfill_start': '2021-11-15', # 回填开始日期 'backfill_end': '2023-10-01' # 回填结束日期 }, reset_dag_run = True, wait_for_completion = True )
修改target_dag的t1执行逻辑:
from datetime import datetime, timedelta def t1(**context): dag_run_conf = context.get('dag_run').conf # 读取回填时间范围 backfill_start = datetime.strptime(dag_run_conf.get('backfill_start'), "%Y-%m-%d") backfill_end = datetime.strptime(dag_run_conf.get('backfill_end'), "%Y-%m-%d") current_date = backfill_start while current_date <= backfill_end: # 此处替换为你原有的单日数据处理逻辑,将current_date作为业务日期传入 process_single_day_business_logic(current_date) current_date += timedelta(days=1)
内容的提问来源于stack exchange,提问作者Tessa Altman
相关产品推荐
相关产品推荐

