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 и завершается успехом в любом другом случае. Но это работает совсем не так как я ожидаю, при первом запуске дага все хорошо, работает корректно, если я меняю данные в текстовом файле, то родительский таск отрабатывает корректно, например запускаю родительский даг заведомо в ошибкой, все отработает корректно дочерний класс будет завершаться с ошибкой, но если я меняю текст, опять родительский отработает корректно, а вот дочерний продолжит падать еще некоторое время, потом возможно будет корректно, но не факт. Если запускать заведомо успешно, ситуация та же, с точностью наоборот. Так же я не понимаю как организовать ожидание.


Ответы (0 шт):