【发布时间】:2018-10-13 19:40:16
【问题描述】:
我希望一个 dag 在另一个 dag 完成后开始。一种解决方案是使用外部传感器功能,您可以在下面找到我的解决方案。我遇到的问题是依赖的 dag 卡在戳,我检查了这个answer 并确保两个 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),
'retries': 1,
'retry_delay': timedelta(minutes=5),
}
schedule = '* * * * *'
dag = DAG('leader_dag', default_args=default_args,catchup=False,
schedule_interval=schedule)
t1 = BashOperator(
task_id='print_date',
bash_command='date',
dag=dag)
依赖的 dag:
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
from airflow.operators.sensors import ExternalTaskSensor
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2018, 10, 8),
'retries': 1,
'retry_delay': timedelta(minutes=5),
}
schedule='* * * * *'
dag = DAG('dependent_dag', default_args=default_args, catchup=False,
schedule_interval=schedule)
wait_for_task = ExternalTaskSensor(task_id = 'wait_for_task',
external_dag_id = 'leader_dag', external_task_id='t1', dag=dag)
t1 = BashOperator(
task_id='print_date',
bash_command='date',
dag=dag)
t1.set_upstream(wait_for_task)
依赖 dag 的日志:
【问题讨论】:
-
感谢您在问题中提供的所有详细信息,尽管将日志作为文本提供会使搜索更加友好。我认为传感器存在根本性的误用,在我的答案末尾显示了修复,但充其量,我不会在 Airflow 中安排比
*/10 * * * *更频繁的任何事情。