【问题标题】:Airflow2 Creating Dag dynamically after function RunAirflow2 在函数运行后动态创建 Dag
【发布时间】:2021-06-27 22:33:27
【问题描述】:

您好,我在这里使用气流是我正在尝试解决的场景 我想在函数运行后动态创建 DAG

try:
    import os
    import sys

    from datetime import timedelta,datetime
    from airflow import DAG

    from airflow.operators.python_operator import PythonOperator
    from airflow.operators.email_operator import EmailOperator
    from airflow.utils.trigger_rule import TriggerRule
    from airflow.utils.task_group import TaskGroup
    import pandas as pd

    print("All Dag modules are ok ......")

except Exception as e:
    print("Error  {} ".format(e))


# ===============================================
default_args = {
    "owner": "airflow",
    "start_date": datetime(2021, 1, 1),
    "retries": 1,
    "retry_delay": timedelta(minutes=1),
    'email': ['shahsoumil519@gmail.com'],
    'email_on_failure': True,
    'email_on_retry': False,
}
dag = DAG(dag_id="project", schedule_interval="@once", default_args=default_args, catchup=False)
# ================================================


class XcomHelper(object):

    def __init__(self, **context):
        self.context = context

    def get(self, key=None):
        """ Get the Value from XCOM"""
        try:
            return self.context.get("ti").xcom_pull(key=key)
        except Exception as e: return "Error"

    def push(self, key=None, value=None):

        """Push the value on session """
        try:
            self.context['ti'].xcom_push(key=key, value=value)
            return True
        except Exception as e: return False



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

    def hello_world_py():
        print('Hello World')

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

    with dag:
        t1 = PythonOperator(task_id=dag_id,python_callable=hello_world_py)

    return dag


def simple_task(**context):

    DATA = ["soumil", "Shah"]
    
    for n in range(1, len(DATA)):
        try:
            dag_id = 'hello_world_{}'.format(str(n))
            print("DAG ID : {} ".format(dag_id))
            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)
        except Exception as e:
            print("Error : {} ".format(e))

with DAG(dag_id="project", schedule_interval="@once", default_args=default_args, catchup=False) as dag:

    simple_task = PythonOperator(task_id="simple_task",
                                 python_callable=simple_task,
                                 provide_context=True)


simple_task


我想根据 DATA 变量的 len 创建这些 dag 该数据来自数据库

我试着调查

任何帮助都会很棒

修改代码:

try:
    import os
    import sys

    from datetime import timedelta, datetime
    from airflow import DAG

    from airflow.operators.python_operator import PythonOperator

    # from airflow.operators.email_operator import EmailOperator
    # from airflow.utils.trigger_rule import TriggerRule
    # from airflow.utils.task_group import TaskGroup
    # import pandas as pd

    print("All Dag modules are ok ......")

except Exception as e:
    print("Error  {} ".format(e))


def create_dag(dag_id, schedule, dag_number, default_args):
    def hello_world_py():
        print('Hello World')

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

    with dag:
        t1 = PythonOperator(task_id=dag_id, python_callable=hello_world_py)

    return dag


def simple_task():
    DATA = ["soumil", "Shah", "Shah2"]

    for n in range(0, len(DATA)):
        try:
            dag_id = 'hello_world_{}'.format(str(n))
            print("DAG ID : {} ".format(dag_id))
            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)
        except Exception as e:
            print("Error : {} ".format(e))

    def trigger_function():
        print("HEREE")
        simple_task()

    with DAG(dag_id="project", schedule_interval="@once", default_args={'owner': 'airflow', 'start_date': datetime(2018, 1, 1)}, catchup=False) as dag:


        trigger_function = PythonOperator(task_id="trigger_function",python_callable=trigger_function,provide_context=True,)


    trigger_function

【问题讨论】:

  • 你需要 dag '项目'吗?您当前的代码正在创建一个生成 dag 的 dag。
  • 好吧,我想基于该函数生成 DAG,一旦我执行该函数,它会说返回 2 项,现在我想做的是创建 2 个 DAG,如果这有意义的话

标签: airflow-scheduler airflow


【解决方案1】:

我从您的代码中删除了几行以保持答案的重点。下面的代码会根据DATA的内容生成像hello_world_0, hello_world_1...这样的DAG。

编辑 - 我使用了气流 v1.10.x,但代码应该适用于 v2.x

建议:

  1. 使任务名称与 DAG 名称不同。
  2. dag_number 变量当前未被使用。可以取下来。

DAG 看起来像这样 -

try:
    import os
    import sys

    from datetime import timedelta, datetime
    from airflow import DAG

    from airflow.operators.python_operator import PythonOperator

    # from airflow.operators.email_operator import EmailOperator
    # from airflow.utils.trigger_rule import TriggerRule
    # from airflow.utils.task_group import TaskGroup
    # import pandas as pd

    print("All Dag modules are ok ......")

except Exception as e:
    print("Error  {} ".format(e))


def create_dag(dag_id, schedule, dag_number, default_args):
    def hello_world_py():
        print('Hello World')

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

    with dag:
        t1 = PythonOperator(task_id=dag_id, python_callable=hello_world_py)

    return dag


def simple_task():
    DATA = ["soumil", "Shah", "Shah2"]

    for n in range(0, len(DATA)):
        try:
            dag_id = 'hello_world_{}'.format(str(n))
            print("DAG ID : {} ".format(dag_id))
            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)
        except Exception as e:
            print("Error : {} ".format(e))


simple_task()

【讨论】:

  • 非常感谢,但是当我手动按下 UI 上的按钮时,我想运行 DAG 创建过程
  • 我已经更新了上面的描述,我只想在从 UI 手动按下时运行,你能帮忙吗
  • SubDag 可能更合适,但用例经过修改。通常不建议使用 DAG 生成另一个 DAG,我会说不可能。
  • 非常感谢这里是我想要实现的用例。有网站 A 现在输入有关任务的信息,现在在气流中添加到数据库中我想根据用户输入的信息创建 Dag 并动态生成这个 dag 问题是当用户添加一个项目时我想为其创建一个 Dag现在我必须重定向到气流仪表板,它将运行脚本并使 dag 你认为这里最好的方式是什么
猜你喜欢
  • 1970-01-01
  • 2014-03-15
  • 2012-07-02
  • 2021-12-27
  • 1970-01-01
  • 1970-01-01
  • 2012-06-16
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多