Airflow多任务分支调度咨询:周一两次运行特殊需求实现
解决方案
核心思路
通过**Airflow变量(Variable)**记录周一当天的运行次数,结合自定义分支逻辑替代单纯的BranchDayOfWeekOperator,区分周一的首次与后续运行;同时在每日任务结束后重置变量,避免影响次日流程。
具体实现步骤
1. 编写分支判断逻辑
自定义Python函数,结合执行日期的星期属性、Airflow变量记录的当日运行次数,决定任务流向:
from airflow.models import Variable import pendulum def determine_task_flow(**context): execution_date = context['execution_date'] # pendulum中周一对应weekday()返回0 is_monday = execution_date.weekday() == 0 if not is_monday: # 非周一直接走task_1→task_2流程 return 'task_1' else: # 读取当日运行计数变量,默认值为0 run_count_key = f"monday_run_count_{execution_date.date()}" run_count = int(Variable.get(run_count_key, default_var=0)) if run_count == 0: # 周一首次运行,更新计数为1,直接执行task_2 Variable.set(run_count_key, 1) return 'task_2' else: # 周一后续运行,走task_1→task_2流程 return 'task_1'
2. 构建DAG与任务依赖
使用PythonBranchOperator实现分支逻辑,同时添加每日重置变量的任务,确保变量不会残留到次日:
from airflow import DAG from airflow.operators.python import PythonOperator, PythonBranchOperator from airflow.utils.dates import days_ago default_args = { 'owner': 'airflow', 'start_date': days_ago(1), } with DAG( 'monday_special_task_flow', default_args=default_args, schedule_interval='@hourly', # 可根据实际需求调整调度频率(如@daily) catchup=False, ) as dag: # 分支判断任务 branch_task = PythonBranchOperator( task_id='branch_task', python_callable=determine_task_flow, provide_context=True, ) # 定义业务任务 task_1 = PythonOperator( task_id='task_1', python_callable=lambda: print("Executing task_1"), ) task_2 = PythonOperator( task_id='task_2', python_callable=lambda: print("Executing task_2"), ) # 重置周一运行计数的任务 def reset_monday_variable(**context): execution_date = context['execution_date'] run_count_key = f"monday_run_count_{execution_date.date()}" try: Variable.delete(run_count_key) except Variable.DoesNotExist: # 变量不存在时忽略异常 pass reset_variable = PythonOperator( task_id='reset_monday_variable', python_callable=reset_monday_variable, provide_context=True, trigger_rule='all_done', # 无论分支走向如何,都执行重置 ) # 设置任务依赖 branch_task >> [task_1, task_2] task_1 >> task_2 task_2 >> reset_variable
3. 关键细节说明
- 变量隔离:用
monday_run_count_YYYY-MM-DD格式命名变量,确保每个周一的计数独立,不会互相干扰。 - 重置时机:每日任务结束后删除当日变量,避免次日读取到旧值;即使变量不存在,捕获异常也不会中断流程。
- 调度适配:无论调度间隔是每小时、每日还是自定义Cron,该逻辑都能准确区分周一的首次与后续运行。
替代方案:基于TaskInstance运行次数判断
如果不想依赖外部变量,可直接查询Airflow内部的TaskInstance运行次数:
from airflow.models import TaskInstance def determine_task_flow(**context): execution_date = context['execution_date'] ti = TaskInstance(task=context['task'], execution_date=execution_date) # 获取当前execution_date对应的任务尝试次数 run_count = ti.get_num_runs() is_monday = execution_date.weekday() == 0 if not is_monday: return 'task_1' else: return 'task_2' if run_count == 1 else 'task_1'
这种方式无需维护外部状态,适合不需要跨任务共享计数的场景。
内容的提问来源于stack exchange,提问作者imarevic
相关产品推荐
相关产品推荐

