Apache Airflow
Quick revision on DAGs, Operators, and TaskFlow orchestration.
The Heart of Orchestration
Apache Airflow is a platform to programmatically author, schedule, and monitor workflows. Workflows are represented as Directed Acyclic Graphs (DAGs).
DAG Orchestration Flow
Below is an animated representation of how a standard ETL DAG executes, respecting dependencies:
Idempotency is Key
Core Concepts
Airflow utilizes several primitive concepts to build complex orchestration layers:
- DAG (Directed Acyclic Graph): The blueprint of your workflow. It defines the execution order.
- Operator: A template for a specific task (e.g.,
PythonOperator,BashOperator,BigQueryInsertJobOperator). - Task: An instantiated Operator inside a DAG.
- Sensor: A special type of operator that waits for an external event to occur before succeeding.
- XCom (Cross-Communication): A mechanism that allows tasks to talk to each other by sharing small amounts of metadata.
Note: Airflow is an orchestrator, not a processing engine. Avoid passing large datasets through XComs. Use XComs to pass URIs or metadata instead!
The TaskFlow API
Modern Airflow (2.0+) strongly recommends using the TaskFlow API which utilizes Python decorators (@dag, @task) to dramatically reduce boilerplate code.
from datetime import datetime
from airflow.decorators import dag, task
@dag(
schedule="@daily",
start_date=datetime(2026, 1, 1),
catchup=False,
tags=["etl", "sales"]
)
def sales_processing_pipeline():
@task
def extract_sales_data():
# logic to extract data
return "gs://bucket/sales_2026_10_10.csv"
@task
def transform_data(file_uri: str):
# logic to transform
return "gs://bucket/sales_clean.csv"
@task
def load_to_warehouse(file_uri: str):
# logic to load into BigQuery
print(f"Loaded {file_uri} successfully!")
# Setting dependencies is as easy as standard Python function calls
raw_file = extract_sales_data()
clean_file = transform_data(raw_file)
load_to_warehouse(clean_file)
# Instantiate the DAG
sales_processing_pipeline()