Airflow: ExternalTaskSensor разное расписание у задач,
Есть два дага Parent и Child, у parent свое расписание предположим '30 * * * *', у child '1 8-17 * * 1-5', child ждет выполнение parent, например 40 минут, если в этот интервал parent завершался с ошибкой, то child тоже падает с ошибкой, в противном случает выполняется следующий таск у дочернего класса. Проблема заключается в том что это не работает даже в самом простом случае, я не пойму как их синхронизировать. Я написал такой код:
Для Parent
import time
from airflow import DAG
from airflow.operators.dummy_operator import DummyOperator
from airflow.operators.python_operator import PythonOperator
from airflow.sensors.external_task_sensor import ExternalTaskSensor, ExternalTaskMarker
start_date = datetime(2021, 3, 1, 20, 36, 0)
class Exept(Exception):
pass
def wait():
time.sleep(3)
with open('etl.txt', 'r') as txt:
line = txt.readline()
if line == 'err':
print(1)
raise Exept
return 'etl success'
with DAG(
dag_id="dag_etl1",
start_date=start_date,
schedule_interval='* * * * *',
tags=['example2'],
) as etl1:
parent_task = ExternalTaskMarker(
task_id="parent_task",
external_dag_id="dag_etl1",
external_task_id="etl_child",
)
wait_timer = PythonOperator(task_id='wait_timer', python_callable=wait)
wait_timer >> parent_task
Для child
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.dummy_operator import DummyOperator
from airflow.operators.python_operator import PythonOperator
from airflow.sensors.external_task_sensor import ExternalTaskSensor, ExternalTaskMarker
from etl_parent import etl1, wait_timer, parent_task
start_date = datetime(2021, 3, 1, 20, 36, 0)
def check():
return 'I succeeded'
with DAG(
dag_id='etl_child',
start_date=start_date,
schedule_interval='* * * * *',
tags = ['testing_child_dag']
) as etl_child:
status = ExternalTaskSensor(
task_id="dag_etl1",
external_dag_id=etl1.dag_id,
external_task_id=parent_task.task_id,
allowed_states=['success'],
mode='reschedule',
execution_delta=timedelta(minutes=1),
timeout=60,
)
task1 = PythonOperator(task_id='task1', python_callable=check)
status >> task1
Как видно, я пытаюсь эмулировать ситуацию когда родительский таск проваливается, если в текстовом файле указан err и завершается успехом в любом другом случае. Но это работает совсем не так как я ожидаю, при первом запуске дага все хорошо, работает корректно, если я меняю данные в текстовом файле, то родительский таск отрабатывает корректно, например запускаю родительский даг заведомо в ошибкой, все отработает корректно дочерний класс будет завершаться с ошибкой, но если я меняю текст, опять родительский отработает корректно, а вот дочерний продолжит падать еще некоторое время, потом возможно будет корректно, но не факт. Если запускать заведомо успешно, ситуация та же, с точностью наоборот. Так же я не понимаю как организовать ожидание.