You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 09:25:26