Условное выполнение задач в AirflowPython

Программы на Python
Anonymous
Условное выполнение задач в Airflow

Сообщение Anonymous »

Я пытаюсь написать группу обеспечения доступности баз данных, которая условно выполняет другую задачу. Упрощенная версия того, с чем я работаю, такова:

Код: Выделить всё

to_be_triggered = EmptyOperator(task_id="to_be_triggered")

@task.branch()
def trigger_dag(**kwargs):
config = kwargs.get("dag_run_config")

if config.get("run_trigger") is True:
return ["to_be_triggered"]
return None

with DAG("example") as dag:
dag_run_config = {
"run_trigger": True
}

t0 = trigger_dag(dag_run_config=dag_run_config)
t1 = EmptyOperator(task_id="end", trigger_rule=TriggerRule.ONE_SUCCESS)

t0 >> t1
Поэтому я хочу условно запустить to_be_triggered, если переменная run_trigger в конфигурации имеет значение True. Я не могу этого сделать, потому что ветка_task_ids должна содержать только допустимые идентификаторы задач, а to_be_triggered по какой-то причине недействителен:

Код: Выделить всё

Following branch {'to_be_triggered'}
Task failed with exception
AirflowException: 'branch_task_ids' must contain only valid task_ids. Invalid tasks found: {'to_be_triggered'}
Судя по данным Google, обычно это происходит потому, что задача находится в группе задач, и ее необходимо указать с помощью идентификатора группы, но у меня здесь нет группы задач. Кто-нибудь знает, задана ли группа задач где-либо неявно или есть ли другая возможная причина того, что to_be_triggered является недействительным?

Вернуться в «Python»