Airflow:如何判断DAG是否为当日最后一次调度运行?
实现Airflow DAG调度触发及当日最后一次调度的判断
要实现这个需求,你可以通过Airflow提供的运行上下文和调度规则计算,分两步完成判断:
1. 判断当前DAG运行是否为调度触发
Airflow的DagRun对象包含运行类型信息,通过run_type字段可区分调度触发和手动触发:
scheduled:调度触发的运行manual:手动触发的运行
在Python任务中,可通过上下文获取dag_run对象进行判断:
from airflow.decorators import task @task def check_and_execute_extra(**context): # 获取当前DAG运行实例 dag_run = context["dag_run"] # 非调度触发直接跳过 if dag_run.run_type != "scheduled": return
2. 判断是否为当日最后一次调度运行
由于你的DAG每日固定运行12次,先明确schedule_interval(比如每2小时一次对应0 */2 * * *),再通过调度规则计算当日最后一次的execution_date,和当前运行的execution_date对比即可。
结合第一步的判断,完整代码示例:
from airflow.decorators import dag, task from airflow.utils.dates import croniter from datetime import datetime, timedelta import pendulum # 定义DAG,假设时区为Asia/Shanghai @dag( schedule_interval="0 */2 * * *", # 每日12次,每2小时运行一次 start_date=datetime(2024, 1, 1, tzinfo=pendulum.timezone("Asia/Shanghai")), catchup=False, timezone="Asia/Shanghai" ) def daily_scheduled_dag(): @task def check_and_execute_extra(**context): dag_run = context["dag_run"] # 第一步:判断是否为调度触发 if dag_run.run_type != "scheduled": print("非调度触发,跳过额外操作") return execution_date = context["execution_date"] dag = context["dag"] day_start = execution_date.replace(hour=0, minute=0, second=0, microsecond=0) day_end = day_start + timedelta(days=1) # 生成当日所有调度的execution_date,取最后一个 cron = croniter(dag.schedule_interval, day_start) last_exec_date = None while True: next_exec = cron.get_next(datetime) # 转换为和execution_date一致的时区 next_exec = pendulum.instance(next_exec).in_timezone(dag.timezone) if next_exec >= day_end: break last_exec_date = next_exec # 第二步:判断是否为当日最后一次调度 if execution_date == last_exec_date: print("执行当日最后一次调度的额外操作") # 这里写入你的额外操作逻辑,比如调用其他任务、执行数据同步等 # your_extra_operation() else: print("非当日最后一次调度,跳过额外操作") # 定义任务依赖 check_and_execute_extra() daily_scheduled_dag()
注意事项
- 时区一致性:确保DAG定义的时区和计算时使用的时区一致,避免因时区差导致判断错误
- 调度规则匹配:
schedule_interval要和实际运行次数对应(每日12次即每2小时一次),如果后续调整调度频率,需同步修改计算逻辑 - 避免依赖数据库查询:直接通过调度规则计算最后一次运行时间比查询元数据库更可靠,不会受未触发调度或延迟运行的影响
内容的提问来源于stack exchange,提问作者SharkyShark
相关产品推荐
相关产品推荐

