Apache Airflow Proficient¶
🔄 Data Engineering · Level 4
When you'd use this
DAGs, operators, scheduling, dependencies and workflow orchestration.
Orchestrate data pipelines as DAGs with scheduling, retries, and dependencies — for reliable recurring batch workflows.
DAG (Directed Acyclic Graph)¶
Define a pipeline as tasks with dependencies that Airflow schedules and retries.
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
default_args = {
"owner": "data-team",
"retries": 3,
"retry_delay": timedelta(minutes=5),
"email_on_failure": True,
"email": ["alerts@company.com"],
}
with DAG(
dag_id="daily_etl_pipeline",
default_args=default_args,
description="Daily user data ETL",
schedule_interval="0 6 * * *", # 6 AM daily
start_date=datetime(2026, 1, 1),
catchup=False,
tags=["etl", "users"],
) as dag:
extract = PythonOperator(
task_id="extract_data",
python_callable=extract_from_source,
)
transform = PythonOperator(
task_id="transform_data",
python_callable=transform_data,
)
validate = PythonOperator(
task_id="validate_data",
python_callable=run_data_quality_checks,
)
load = PythonOperator(
task_id="load_to_warehouse",
python_callable=load_to_bigquery,
)
notify = BashOperator(
task_id="send_notification",
bash_command='echo "ETL complete" | mail -s "Daily ETL" team@company.com',
)
# Define dependencies
extract >> transform >> validate >> load >> notify
TaskFlow API (Python-native, modern)¶
Write tasks as decorated Python functions with automatic data passing.
from airflow.decorators import dag, task
from datetime import datetime
@dag(
schedule_interval="@daily",
start_date=datetime(2026, 1, 1),
catchup=False,
)
def etl_pipeline():
@task
def extract() -> dict:
"""Extract data from source."""
import pandas as pd
df = pd.read_csv("s3://bucket/raw/users.csv")
return {"rows": len(df), "path": "s3://bucket/raw/users.csv"}
@task
def transform(extract_result: dict) -> str:
"""Transform and save to staging."""
import pandas as pd
df = pd.read_csv(extract_result["path"])
df["email"] = df["email"].str.lower()
output_path = "s3://bucket/staging/users_clean.parquet"
df.to_parquet(output_path)
return output_path
@task
def load(path: str):
"""Load to data warehouse."""
import pandas as pd
df = pd.read_parquet(path)
# load to BigQuery/Redshift...
print(f"Loaded {len(df)} rows")
# Automatic dependency inference
data = extract()
cleaned = transform(data)
load(cleaned)
etl_pipeline() # register the DAG
Practice Exercises¶
- Build a DAG with extract → transform → load → validate → notify.
- Add branching — different transforms based on data characteristics.
- Implement retry and alerting — 3 retries with exponential backoff, email on failure.
- Use XCom to pass data between tasks.
- Schedule backfills — process historical data for a date range.
💬 Discussion
Have a question about this topic? Found an error? Share your thoughts below.