【问题标题】:How to dynamically delete DAGs?如何动态删除 DAG?
【发布时间】:2021-02-19 20:37:13
【问题描述】:

我正在动态创建 DAG,按照 Dynamically Generating DAGs in Airflow 中的说明,通过变量 k 修改要创建的 dag 数量:

from datetime import datetime
from airflow import DAG
from airflow.operators.python_operator import PythonOperator

def create_dag(dag_id,
               schedule,
               dag_number,
               default_args):

    def hello_world_py(*args):
        print('Hello World')
        print('This is DAG: {}'.format(str(dag_number)))

    dag = DAG(dag_id,
              schedule_interval=schedule,
              default_args=default_args)

    with dag:
        t1 = PythonOperator(
            task_id='hello_world',
            python_callable=hello_world_py,
            dag_number=dag_number)

    return dag


# build k dags
k = 5
for n in range(1, k + 1):
    dag_id = 'hello_world_{}'.format(str(n))

    default_args = {'owner': 'airflow',
                    'start_date': datetime(2018, 1, 1)
                    }

    schedule = '@daily'

    dag_number = n

    globals()[dag_id] = create_dag(dag_id,
                                  schedule,
                                  dag_number,
                                  default_args)

我可以使用 UI 和 CLI 检查创建的 DAG。两者同步:

> airflow dags list
dag_id        | filepath       | owner   | paused
==============+================+=========+=======
hello_world_1 | hello_world.py | airflow | True
hello_world_2 | hello_world.py | airflow | True
hello_world_3 | hello_world.py | airflow | True
hello_world_4 | hello_world.py | airflow | True
hello_world_5 | hello_world.py | airflow | True

现在,如果我将 k 减少到 3,CLI 将仅按预期列出 3 个 dag。但是 UI 一直显示 5 个 dag。

如何使 UI 与要创建的 dag 数量保持同步?如何在 python 中以编程方式删除 DAG?我想像创建 DAG 一样轻松地删除它们。

【问题讨论】:

  • dag 可能会被调度程序再次选中,因为 .py 文件存在并正在处理中

标签: python airflow airflow-scheduler


【解决方案1】:

这是一种动态删除 dag 的方法。

from airflow.api.common.experimental.delete_dag import delete_dag
from airflow.utils.session import provide_session
from airflow.models import DagModel

enabled_dags = [] # dynamic enabled_dags

@provide_session
def get_all_dag_ids(session=None):
    all_objs = session.query(DagModel).all()
    return [i.dag_id for i in all_objs]

all_dag_ids = get_all_dag_ids() # all dag in database

for k in [gk for gk in all_dag_ids if not in enabled_dags]:
    delete_dag(k)
    del globals()[k]

airflow 只是加载 dag 文件总是添加到数据库中。如果 delete dag with dag_id query all dag_id from db, 然后删除。

希望这有帮助。

【讨论】:

    猜你喜欢
    • 2017-03-31
    • 1970-01-01
    • 1970-01-01
    • 2014-03-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多