
在 Airflow 中,可通过组合定时调度、任务依赖与 TimeSensorAsync(或 TimeDeltaSensorAsync)实现:让 T08:30 任务仅在当日 T05:30 任务成功完成 且 到达指定时间后才触发,兼顾时序控制与状态依赖。
在 airflow 中,可通过组合定时调度、任务依赖与 timesensorasync(或 timedeltasensorasync)实现:让 t08:30 任务仅在当日 t05:30 任务成功完成 *且* 到达指定时间后才触发,兼顾时序控制与状态依赖。
Airflow 的 DAG 级调度(schedule_interval 或 schedule)决定了整个工作流的触发节奏,而任务级依赖(>>)确保执行顺序与状态约束。要同时满足「每日固定时间触发」和「仅当同日前置任务成功才继续」两个条件,不能仅靠 DAG 调度或单纯上下游依赖——因为默认调度下,T08:30 任务若独立设为每天 08:30 触发,将无视 T05:30 的当日执行结果。
✅ 正确解法是:将整个 DAG 设为每日 05:30 触发(即以 T05:30 为起点),并在其后插入一个时间传感器,使 T08:30 任务延迟至当天 08:30 才运行,且仅在 T05:30 成功的前提下推进。
以下是推荐实现(基于 Airflow 2.6+,使用 TimeSensorAsync 提升效率):
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.sensors.time import TimeSensorAsync # 异步版,不阻塞 Worker
from datetime import datetime, time
default_args = {
"owner": "airflow",
"retries": 1,
}
with DAG(
dag_id="t0530_then_t0830",
default_args=default_args,
schedule="0 5 * * *", # 每日 05:30 UTC 触发整个 DAG(注意:Airflow 默认用 UTC)
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["timing", "dependency"],
) as dag:
def task_t0530():
print("Running T05:30 job — e.g., data ingestion")
return "success"
t0530 = PythonOperator(
task_id="T05:30",
python_callable=task_t0530,
)
# 等待到当日 08:30(UTC),但仅在 t0530 成功后才开始等待
wait_until_0830 = TimeSensorAsync(
task_id="wait_until_0830",
target_time=time(8, 30), # 注意:target_time 是本地时间(即 DAG timezone 下的时间)
mode="reschedule", # 推荐:释放 worker slot,避免长时间占用
)
def task_t0830():
print("Running T08:30 job — e.g., downstream analytics")
return "completed"
t0830 = PythonOperator(
task_id="T08:30",
python_callable=task_t0830,
)
# 关键依赖链:T05:30 → 等待 → T08:30
t0530 >> wait_until_0830 >> t0830
? 关键说明与注意事项:
- ✅ 时区对齐:
TimeSensorAsync.target_time解析为 DAG 所配置的timezone(如未显式设置,默认为UTC)。若需按北京时间(CST/UTC+8)执行,请在 DAG 中声明:from airflow.utils.timezone import utc from pendulum import timezone tz_beijing = timezone("Asia/Shanghai") # 然后在 DAG 参数中加入:timezone=tz_beijing - ⚠️ 避免
TimeSensor(同步版):它会持续轮询并占用一个 Worker Slot,尤其在长等待(如 3 小时)时资源浪费严重;TimeSensorAsync基于异步事件驱动,更高效可靠。 - ?
catchup=False必须启用:防止历史调度堆积导致多个“当日”实例并发,破坏「仅执行一次」语义。 - ? 替代方案(高级场景):若需跨多日依赖(如“T08:30 必须等 T05:30 在 同一逻辑日期 成功”,而 DAG 运行时间因延迟偏移),可结合
ExternalTaskSensor+ 自定义execution_date_fn,但本例中单一 DAG + 时间传感器已足够简洁健壮。
综上,该方案完全符合 Airflow 的设计哲学——用原生 Operator 表达业务意图,无需外部变量、数据库标记或 hack 式轮询,是处理“定时 + 状态门控”类需求的地道(idiomatic)实践。










