【发布时间】:2018-04-24 13:33:39
【问题描述】:
我正在尝试在测试环境中测试具有多个任务的 dag。我能够测试与 dag 关联的单个任务,但我想在 dag 中创建多个任务并启动第一个任务。 为了在 dag 中测试一项任务,我正在使用
task1.run()
正在执行。但是,当我在一个 dag 的下游一个接一个地执行许多任务时,同样的方法不起作用。
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2015, 6, 1),
'email': ['airflow@example.com'],
'email_on_failure': False,
'email_on_retry': False,
'retries': 1,
'retry_delay': timedelta(minutes=5),
# 'queue': 'bash_queue',
# 'pool': 'backfill',
# 'priority_weight': 10,
# 'end_date': datetime(2016, 1, 1),
}
dag = DAG('tutorial', default_args=default_args)
# t1, t2 and t3 are examples of tasks created by instantiating operators
t1 = BashOperator(
task_id='print_date',
bash_command='date',
dag=dag)
t2 = BashOperator(
task_id='sleep',
bash_command='sleep 5',
retries=3,
dag=dag)
t2.set_upstream(t1)
t1.run() # It is executing just first task.
为了运行第二个任务,我必须使用 t2.run() 来运行,因为我在设计 DAG 时不希望这样做。如何做到这一点?
【问题讨论】:
-
我正在尝试测试 dag 的整个运行过程。如果我们将上面的代码放在测试语句中,它就不起作用。底线是我愿意使用 pytest 测试由许多依赖任务组成的整个 dag 运行。
-
那个指向 Dag 的单元测试的链接只涉及一个任务。我想测试整个 dag 而不是每个 dag 只测试一个任务。希望这能消除混乱。
-
到目前为止你尝试过什么?测试代码是什么样子的?
dag.run()呢? -
我不知道为什么它没有点击。简单的 dag.run() 有效。谢谢
标签: python unit-testing pytest airflow