Airflow not picking up dags from a helper class

Viewed 324

I tried to experiment with a helper class that can create a dag with the parameters passed in . but when I try to import the class in my dag file, airflow doesn't pick it up.

Here's my helper class:

from airflow.models.dag import ScheduleInterval
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta, date

class dagClass:
    def __init__(self) -> None:
        self.dag = None

        self.args={
        'owner': 'airflow',
        'depends_on_past': False,
        'start_date': datetime.now()-timedelta(days=1),
        'email': ['airflow@airflow.com'],
        'email_on_failure': False,
        'email_on_retry': False,
        'retries': 1,
        'retry_delay': timedelta(minutes=1),
        }
    def create_dag(self,dagName):
        print(type(self.args))
        self.dag=DAG(dagName, 
        default_args=self.args,
        schedule_interval="39 8 * * *")
        return self.dag
    def get_dag(self, dagName):
        dag = self.create_dag(dagName)
        print(type(dag))
        return dag

    def add_task(self,dag,taskID, scriptPath):
        return BashOperator(
            task_id = taskID,
            bash_command = "python3 {}".format(scriptPath),
            dag = dag
        )

and here's how I am importing it

from dagClass import dagClass

dagObj = dagClass()

dag = dagObj.get_dag("testClass21")

t1 = dagObj.add_task(dag,"t1","test.py")


print(type(dag)) #<class 'airflow.models.dag.DAG'>
print(type(t1)) # <class 'airflow.operators.bash.BashOperator'>

t1

anybody has any idea on what's wrong? airflow is picking up the other dags I created so it isn't a problem with the folder.

2 Answers

Tried out @NicoE 's suggestion and it worked, it seems like airflow scans for keywords like DAG or airflow in the py files and if it finds them marks them as dags. So you can either add those imports to make airflow recognize them or set dag_discovery_safe_mode to False in airflow.cfg to disable this kind of checking

Adding more context to my previous comment. From Loading DAGs section of the docs:

When searching for DAGs inside the DAG_FOLDER, Airflow only considers Python files that contain the strings airflow and dag (case-insensitively) as an optimization. To consider all Python files instead, disable the DAG_DISCOVERY_SAFE_MODE configuration flag.

So, in order to achieve that Airflow discover the DAGs, it's neccessary that the DAG objects are declared as top level code of the .py file (you have done that already), and that it contains airflow and DAG strings somewhere within. This behaviour can be changed by setting dag_discovery_safe_mode to False.

Note: I don't think that the imports are mandatory, just the presence of the terms as a string, but I haven't tried it myself.

Related