Airflow不同调度周期DAG协同调度方案合理性咨询
问题描述
我们有多个执行各类数据处理任务的DAG,随着系统扩容,会有更多来自不同内部团队的DAG加入,部分DAG会依赖其他团队的DAG及数据。我们计划通过一个“主调度DAG”,使用TriggerDagRunOperator来协调所有DAG间依赖,示例代码如下:
dag_1 = TriggerDagRunOperator(trigger_dag_id = "dag_1_id", ...) dag_2 = TriggerDagRunOperator(trigger_dag_id = "dag_2_id", ...) dag_3 = TriggerDagRunOperator(trigger_dag_id = "dag_3_id", ...) dag_4 = TriggerDagRunOperator(trigger_dag_id = "dag_4_id", ...) dag_5 = TriggerDagRunOperator(trigger_dag_id = "dag_5_id", ...) # dag_1无依赖,也没有任务依赖它,很简单! dag_1 # 开发dag_3的团队依赖dag_2的输出 dag_2 >> dag_3 # dag_5依赖dag_2和dag_4 # 特殊情况:dag_2只需每日执行一次,但开发dag_5的团队希望它每20分钟运行一次——因为它还有其他更新更频繁的外部依赖,和dag_2的日频输出无关 dag_5 << [dag_2, dag_4]
当前面临调度适配问题:部分DAG仅需每日执行一次,但部分DAG需更频繁执行(如dag_5需每20分钟执行一次,且依赖每日执行的dag_2)。原本设想是移除子DAG的独立调度规则,通过主调度DAG中的时间传感器触发执行,但担忧主调度DAG需按最频繁子DAG的周期运行,觉得该方案不合理,希望验证该方案的合理性,或获取更优实现方案(使用Airflow 2.2.3及Google Cloud Composer)。
可行解决方案
方案一:保留子DAG独立调度,用外部任务传感器替代主调度硬触发
- 对于日频的dag_2,保留自身
schedule_interval="@daily"的调度规则,无需主DAG触发。 - 对于高频的dag_5,设置
schedule_interval="*/20 * * * *",并在dag_5中添加ExternalTaskSensor,监听dag_2当日的成功运行实例,确保每次dag_5运行时能获取到dag_2的最新日更数据:from airflow.sensors.external_task import ExternalTaskSensor wait_for_dag2 = ExternalTaskSensor( task_id='wait_for_dag2', external_dag_id='dag_2_id', external_task_id=None, # 监听dag_2整个DAG的完成 execution_date_fn=lambda dt: dt.floor('D'), # 匹配当日的dag_2执行实例 mode='reschedule', # 节省资源,定期检查状态 timeout=3600*24, # 最长等待1天,避免无限阻塞 ) # dag_5的核心任务依赖该传感器和其他外部依赖任务 wait_for_dag2 >> dag_5_main_task - 优势:各DAG的调度周期由所属团队自主维护,无需集中式主调度DAG,避免高频运行的资源浪费;依赖关系由消费方(dag_5)主动监听,符合分布式团队的权责划分。
方案二:拆分主调度DAG逻辑,按需触发不同周期的子DAG
如果坚持使用主调度DAG模式,可拆分触发逻辑避免无意义的高频运行:
- 将主调度DAG拆分为两个独立的主DAG:
- 日频主DAG:
schedule_interval="@daily",负责触发dag_1、dag_2、dag_3这些日频任务,并维护它们的依赖关系(如dag_2 >> dag_3)。 - 高频主DAG:
schedule_interval="*/20 * * * *",负责触发dag_4、dag_5,并在触发dag_5前通过ExternalTaskSensor确认当日dag_2已成功运行。
- 日频主DAG:
- 或者在同一个主DAG中,用
BranchPythonOperator结合时间判断,决定每个调度周期需要触发的子DAG:from airflow.operators.python import BranchPythonOperator def decide_tasks_to_trigger(**context): execution_date = context['execution_date'] tasks_to_trigger = [] # 每日0点触发日频任务 if execution_date.minute == 0 and execution_date.hour == 0: tasks_to_trigger.extend(['trigger_dag1', 'trigger_dag2', 'trigger_dag3']) # 每20分钟触发高频任务 tasks_to_trigger.extend(['trigger_dag4', 'trigger_dag5']) return tasks_to_trigger branch_task = BranchPythonOperator( task_id='decide_tasks', python_callable=decide_tasks_to_trigger, provide_context=True, ) # 定义各TriggerDagRunOperator任务 trigger_dag1 = TriggerDagRunOperator(task_id='trigger_dag1', trigger_dag_id='dag_1_id', ...) trigger_dag2 = TriggerDagRunOperator(task_id='trigger_dag2', trigger_dag_id='dag_2_id', ...) trigger_dag3 = TriggerDagRunOperator(task_id='trigger_dag3', trigger_dag_id='dag_3_id', ...) trigger_dag4 = TriggerDagRunOperator(task_id='trigger_dag4', trigger_dag_id='dag_4_id', ...) trigger_dag5 = TriggerDagRunOperator(task_id='trigger_dag5', trigger_dag_id='dag_5_id', ...) # 设置依赖:dag_2执行完再触发dag_3;触发dag_5前先确认dag_2当日已完成 wait_for_dag2 = ExternalTaskSensor( task_id='wait_for_dag2_for_dag5', external_dag_id='dag_2_id', execution_date_fn=lambda dt: dt.floor('D'), mode='reschedule', ) branch_task >> [trigger_dag1, trigger_dag2, trigger_dag4] trigger_dag2 >> trigger_dag3 [trigger_dag4, wait_for_dag2] >> trigger_dag5 - 优势:主DAG按高频周期运行,但仅在必要时触发日频任务,避免重复触发dag_2;同时统一管理依赖关系,适合需要集中调度管控的场景。
方案三:使用Airflow Dataset功能(Airflow 2.2+支持)
Airflow 2.2及以上版本支持Dataset功能,可通过数据资产的更新来触发DAG,更贴合数据依赖的本质:
- 配置dag_2的输出为一个Dataset:
from airflow import Dataset dag2_output = Dataset("gs://your-bucket/dag2-output/") with DAG( dag_id="dag_2_id", schedule_interval="@daily", ... ) as dag2: # dag_2的最后一个任务产出该Dataset final_task = BashOperator( task_id="final_task", bash_command="gsutil cp ... gs://your-bucket/dag2-output/", outlets=[dag2_output], # 标记该任务产出Dataset ) - 配置dag_5的调度依赖于该Dataset和自身的时间周期:
with DAG( dag_id="dag_5_id", schedule=["*/20 * * * *", dag2_output], # 每20分钟触发一次,或Dataset更新时立即触发 ... ) as dag5: # dag_5的任务逻辑 ... - 优势:无需依赖传感器或主调度DAG,通过数据本身的更新驱动依赖,逻辑更直观;dag_5会按周期运行,同时如果当日dag_2刚完成数据更新,也会立即触发一次,确保数据新鲜度。
原方案合理性分析
你提到的“移除子DAG独立调度,主调度DAG按最频繁周期运行”的方案存在明显不合理性:
- 资源浪费:主DAG每20分钟运行一次,大部分时间只是重复触发无需频繁执行的dag_2,造成不必要的调度开销。
- 数据风险:dag_2每日仅需执行一次,但主DAG频繁触发会导致dag_2重复运行,可能造成数据重复生成、计算资源浪费,甚至数据一致性问题。
- 维护复杂:所有DAG的调度规则都集中在主DAG中,后续新增或修改DAG周期时,都需要修改主DAG,不符合团队协作的独立性要求。
因此该方案不推荐,更建议根据团队管控需求选择方案一或方案三。
内容的提问来源于stack exchange,提问作者FrustratedWithFormsDesigner
相关产品推荐
相关产品推荐

