Last modified: Oct 10, 2026

Write Your First Airflow DAG: Guide

Apache Airflow is a workflow orchestration platform. It lets you define data pipelines as Python code. The core concept is the DAG.

Writing your first DAG can feel intimidating. There are operators, dependencies, schedules, and a web UI to learn. But the basics are simpler than they look.

This guide walks you through building a DAG from scratch. You will learn the structure, the key classes, and how to run and debug your pipeline.

What Is a DAG?

A DAG stands for Directed Acyclic Graph. It is a collection of tasks with defined order and dependencies.

Directed means tasks flow in one direction. Acyclic means there are no loops. A task cannot depend on itself.

Each task is a node in the graph. Edges define what runs before what. Airflow uses this structure to schedule and monitor every step.

If you have not installed Airflow yet, start with our guide on how to install Apache Airflow in Python. It covers the virtual environment, the constraints file, and the database setup.

The Core Building Blocks

Every DAG uses a small set of classes. Understanding them makes the rest straightforward.

The DAG() class is the container. It holds your tasks and defines the schedule. You give it an ID, a start date, and a schedule interval.

The PythonOperator() class wraps a Python function as a task. This is the most common operator for beginners.

The BashOperator() class runs shell commands. Use it when your logic lives in a script.

The EmptyOperator() class does nothing. It is useful as a starting or ending marker in a pipeline.

Create Your DAG Folder

Airflow scans a folder called dags inside your AIRFLOW_HOME directory. Any Python file in that folder is parsed as a DAG.

The default location is ~/airflow/dags. Create it if it does not exist.

 mkdir -p ~/airflow/dags cd ~/airflow/dags 

Every file you place here is scanned by the scheduler. Keep each DAG in its own file for clarity.

Write a Minimal DAG

Start with the smallest possible DAG. One task that prints a message.

 from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator def greet(): print("Hello from my first DAG!") # Define the DAG with DAG( dag_id='my_first_dag', start_date=datetime(2024, 5, 1), schedule_interval='@daily', catchup=False, ) as dag: # Define a single task greet_task = PythonOperator( task_id='greet_task', python_callable=greet, ) 

Save the file as my_first_dag.py in the dags folder. The scheduler picks it up within a minute.

Refresh the web UI at http://localhost:8080. Your DAG appears in the list.

Understanding the Parameters

The dag_id must be unique across all DAGs. It is the identifier you see in the UI.

The start_date tells Airflow when the DAG becomes active. Airflow creates task runs from this date forward. Pick a past date, not a future one.

The schedule_interval defines how often the DAG runs. You can use a cron expression, a timedelta, or a preset like @daily.

The catchup flag controls backfills. When set to False, Airflow only runs the most recent interval. When True, it runs every missed interval from the start date.

Warning: Setting catchup to True with an old start date triggers thousands of runs. Set it to False for new DAGs until you understand backfills.

Add Multiple Tasks with Dependencies

Real pipelines have several steps. You chain tasks with the >> operator or the set_downstream() method.

 from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.empty import EmptyOperator def extract(): print("Extracting data...") return "raw_data" def transform(): print("Transforming data...") def load(): print("Loading data...") default_args = { 'owner': 'admin', 'retries': 2, 'retry_delay': timedelta(minutes=5), } with DAG( dag_id='etl_pipeline', default_args=default_args, start_date=datetime(2024, 5, 1), schedule_interval='@daily', catchup=False, tags=['etl', 'example'], ) as dag: start = EmptyOperator(task_id='start') end = EmptyOperator(task_id='end') extract_task = PythonOperator( task_id='extract', python_callable=extract, ) transform_task = PythonOperator( task_id='transform', python_callable=transform, ) load_task = PythonOperator( task_id='load', python_callable=load, ) # Set the order of execution start >> extract_task >> transform_task >> load_task >> end 

The chain operator >> sets upstream and downstream relationships. Airflow runs tasks in order and only proceeds when the previous task succeeds.

Best practice: Use EmptyOperator as start and end markers. It keeps the graph clean and makes it easy to add branches later.

Pass Data Between Tasks

Tasks are isolated. They cannot share variables directly. Airflow provides XComs for passing small values.

Any value returned from a Python callable is stored as an XCom. Downstream tasks pull it with the xcom_pull() method.

 def extract(**context): data = {"rows": 100, "source": "database"} return data def transform(**context): ti = context['ti'] data = ti.xcom_pull(task_ids='extract') print(f"Received: {data}") return {"rows": data['rows'], "cleaned": True} extract_task = PythonOperator( task_id='extract', python_callable=extract, ) transform_task = PythonOperator( task_id='transform', python_callable=transform, ) extract_task >> transform_task 

XComs are stored in the metadata database. They work well for small values like row counts or file paths.

Important: Do not pass large DataFrames through XComs. Store them in a file or object storage and pass the path instead.

Run and Test the DAG

The scheduler triggers DAGs automatically. But during development, you want manual control.

Use the CLI to list DAGs and trigger a run.

 airflow dags list airflow dags trigger etl_pipeline 
 dag_id | fileloc | owners | is_paused ==============+===============================+========+========== etl_pipeline | /home/user/airflow/dags/... | admin | False my_first_dag | /home/user/airflow/dags/... | admin | False 

You can also test a single task without running the whole DAG. This is faster for debugging.

 airflow tasks test etl_pipeline extract 2024-05-01 
 [2024-05-01 10:15:00,000] {python.py:177} INFO - Extracting data... [2024-05-01 10:15:00,100] {taskinstance.py:1234} INFO - Marking task as SUCCESS 

The tasks test command runs the task in isolation. It does not write to the metadata database and does not affect scheduling.

Common Errors and How to Fix Them

DAG does not appear in the UI: The scheduler has not parsed the file. Check for Python syntax errors. Run airflow dags list to see the parsing error.

Task fails with a missing argument: Python callables in operators cannot take arbitrary arguments unless you pass them explicitly. Use the op_args or op_kwargs parameters.

 def process(date, region): print(f"Processing {region} for {date}") task = PythonOperator( task_id='process_task', python_callable=process, op_kwargs={'date': '2024-05-01', 'region': 'us-east'}, ) 

ImportError inside the DAG: The package is not installed in the Airflow environment. Activate the virtual environment and install it there.

DAG is paused: New DAGs are paused by default. Toggle the switch in the UI or run airflow dags unpause etl_pipeline.

Schedule interval not firing: The start_date must be in the past. Also check that the DAG is unpaused and the scheduler is running.

Heavy work in DAG parse time: Do not run database queries or API calls at the top level of a DAG file. Put that logic inside tasks. The scheduler parses files every 30 seconds.

Best Practices for Writing DAGs

Keep DAG files small. Put complex logic in separate Python modules and import them.

Use default_args for shared settings like retries and owner. This avoids repetition across tasks.

Give tasks clear, descriptive IDs. The UI shows them, and clear names make debugging easier.

Set retries and retry delays. Transient failures are common in data pipelines.

Add tags to categorize DAGs. They help filter the UI when you have dozens of workflows.

Use the catchup=False flag for new DAGs. Enable backfills only when you need them.

Test tasks individually with airflow tasks test before running the whole DAG.

Add notifications for failures. You can send HTML emails with SendGrid and Python from an on_failure_callback to alert your team instantly.

If your pipeline generates reports, produce them inside a task and email them as attachments. A common pattern is to generate PDFs with ReportLab, then attach the result to a SendGrid message.

Conclusion

Writing your first Airflow DAG is simpler than it looks. Define a DAG() object. Add tasks with operators like PythonOperator(). Chain them with the >> operator. Then let the scheduler handle the rest.

Start with a single task. Add dependencies and XComs as your pipeline grows. Test tasks in isolation with the CLI before running the full DAG.

Follow the examples in this guide and you will have a working pipeline in minutes. From there, explore sensors, branching, and custom operators to build more complex workflows. Happy orchestrating!