【发布时间】:2019-06-04 21:40:33
【问题描述】:
我有多个任务正在相互传递一些数据对象。在某些任务中,如果不满足某些条件,我将引发异常。这导致该任务的失败。当触发下一次 DAG 运行时,已经成功的任务再次运行。我正在寻找一些方法来避免运行之前成功的任务,并在下一次 DAG 运行中从失败的任务中恢复 DAG 运行。
【问题讨论】:
-
下一次dag运行时如何触发“已经成功的任务”? Dag 不是在每次运行时都会触发自己的一组任务吗?
我有多个任务正在相互传递一些数据对象。在某些任务中,如果不满足某些条件,我将引发异常。这导致该任务的失败。当触发下一次 DAG 运行时,已经成功的任务再次运行。我正在寻找一些方法来避免运行之前成功的任务,并在下一次 DAG 运行中从失败的任务中恢复 DAG 运行。
【问题讨论】:
如前所述,每个 DAG 都有一组每次运行都会执行的任务。为了避免运行以前成功的任务,您可以通过Airflow XCOMs 或Airflow Variables 对外部变量进行检查,也可以查询元数据库以了解以前运行的状态。您还可以将变量存储在 Redis 或类似的外部数据库中。
使用该变量,您可以跳过任务的执行并直接将任务标记为成功,直到它到达要完成的任务。
当然,如果 DAG 运行时间可能重叠,您需要注意任何潜在的竞争条件。
def task_1( **kwargs ):
if external_variable:
pass
else:
perform_task()
return True
【讨论】: