【发布时间】:2018-08-03 03:07:39
【问题描述】:
在我的 DAG 中,我有一些只应在星期六运行的任务。因此,我使用 BranchPythonOperator 在星期六的任务和 DummyTask 之间进行分支。之后,我加入了两个分支并想运行其他任务。
工作流程如下所示:
这里我将dummy3的触发规则设置为'one_success',一切正常。
我遇到的问题是 BranchPythonOperator 上游出现故障时:
BranchPythonOperator 和分支正确的状态为'upstream_failed',但加入分支的任务变为'skipped',因此整个工作流显示'success'。
我尝试使用'all_success' 作为触发规则,然后如果出现故障则整个工作流程都会失败,但如果没有任何故障,则跳过 dummy3。
我也试过'all_done'作为触发规则,如果没有失败,它会正常工作,但如果有失败dummy3仍然会被执行。
我的测试代码如下所示:
from datetime import datetime, date
from airflow import DAG
from airflow.operators.python_operator import BranchPythonOperator, PythonOperator
from airflow.operators.dummy_operator import DummyOperator
dag = DAG('test_branches',
description='Test branches',
catchup=False,
schedule_interval='0 0 * * *',
start_date=datetime(2018, 8, 1))
def python1():
raise Exception('Test failure')
# print 'Test success'
dummy1 = PythonOperator(
task_id='python1',
python_callable=python1,
dag=dag
)
dummy2 = DummyOperator(
task_id='dummy2',
dag=dag
)
dummy3 = DummyOperator(
task_id='dummy3',
dag=dag,
trigger_rule='one_success'
)
def is_saturday():
if date.today().weekday() == 6:
return 'dummy2'
else:
return 'today_is_not_saturday'
branch_on_saturday = BranchPythonOperator(
task_id='branch_on_saturday',
python_callable=is_saturday,
dag=dag)
not_saturday = DummyOperator(
task_id='today_is_not_saturday',
dag=dag
)
dummy1 >> branch_on_saturday >> dummy2 >> dummy3
branch_on_saturday >> not_saturday >> dummy3
编辑
我刚刚想出了一个丑陋的解决方法:
dummy4 代表我实际需要运行的任务,dummy5 只是一个假人。
dummy3 仍然有触发规则'one_success'。
现在 dummy3 和 dummy4 在上游没有失败的情况下运行,如果当天不是星期六 dummy5 '运行',如果当天是星期六则跳过,这意味着 DAG 在两种情况下都被标记为成功。
如果上游出现故障,则跳过 dummy3 和 dummy4,将 dummy5 标记为'upstream_failed',并将 DAG 标记为失败。
这种解决方法可以让我的 DAG 按我的意愿运行,但我仍然更喜欢没有一些 hacky 解决方法的解决方案。
【问题讨论】: