Data Engineering Path · Airflow
Real-World Pipeline Example
🏗️ Complete Production Pipeline: E-Commerce Analytics
Architecture Overview
graph LR
subgraph "Source Systems"
API["🌐 REST API"]
DB["🗄️ PostgreSQL"]
S3["📦 S3 Raw Data"]
end
subgraph "Airflow Pipeline"
E1["Extract API"] --> T["Transform<br/>(Spark)"]
E2["Extract DB"] --> T
E3["Wait S3 File"] --> T
T --> Q["Quality Check"]
Q --> L["Load to<br/>Snowflake"]
L --> R["Refresh<br/>dbt Models"]
R --> N["Notify<br/>Slack"]
end
subgraph "Target Systems"
SNOW["❄️ Snowflake"]
DASH["📊 Dashboard"]
end
API --> E1
DB --> E2
S3 --> E3
L --> SNOW
R --> DASH
style E1 fill:#4CAF50,stroke:#388E3C,color:#fff
style E2 fill:#4CAF50,stroke:#388E3C,color:#fff
style E3 fill:#9C27B0,stroke:#7B1FA2,color:#fff
style T fill:#FF9800,stroke:#F57C00,color:#fff
style Q fill:#F44336,stroke:#D32F2F,color:#fff
style L fill:#2196F3,stroke:#1976D2,color:#fff
style R fill:#00BCD4,stroke:#0097A7,color:#fff
style N fill:#607D8B,stroke:#455A64,color:#fff
The Complete DAG
from airflow.sdk import dag, task
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.providers.standard.operators.bash import BashOperator
from airflow.operators.empty import EmptyOperator
from datetime import datetime, timedelta
default_args = {
"owner": "data-engineering",
"retries": 3,
"retry_delay": timedelta(minutes=5),
"email_on_failure": True,
"email": ["data-alerts@company.com"],
}
@dag(
dag_id="ecommerce_analytics_v2",
default_args=default_args,
schedule="0 7 * * *",
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["production", "ecommerce", "analytics"],
max_active_runs=1,
doc_md="""
## E-Commerce Analytics Pipeline
This DAG runs daily at 7:00 AM UTC and:
1. Extracts data from 3 sources (API, DB, S3)
2. Transforms with Apache Spark
3. Runs data quality checks
4. Loads to Snowflake
5. Refreshes dbt models
6. Sends Slack notification
""",
)
def ecommerce_analytics():
start = EmptyOperator(task_id="pipeline_start")
wait_for_file = S3KeySensor(
task_id="wait_for_clickstream",
bucket_name="raw-data-prod",
bucket_key="clickstream/{{ ds }}/events.parquet",
deferrable=True,
timeout=7200,
)
@task()
def extract_orders() -> dict:
from airflow.providers.postgres.hooks.postgres import PostgresHook
hook = PostgresHook("orders_db")
records = hook.get_records(
"SELECT * FROM orders WHERE DATE(created_at) = %s",
parameters=[("{{ ds }}",)],
)
return {"count": len(records), "path": f"s3://staging/orders/{{ ds }}"}
@task()
def extract_products() -> dict:
import requests
response = requests.get("https://api.company.com/products")
return {"count": len(response.json()), "data": response.json()}
@task()
def run_spark_transform(orders: dict, products: dict):
from airflow.providers.apache.spark.hooks.spark_submit import SparkSubmitHook
hook = SparkSubmitHook(conn_id="spark_cluster")
hook.submit(application="/spark/jobs/ecommerce_transform.py")
return {"status": "completed"}
@task()
def data_quality_check(transform_result: dict):
from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook
hook = SnowflakeHook("snowflake_prod")
row_count = hook.get_first("SELECT COUNT(*) FROM analytics.daily_metrics")[0]
assert row_count > 0, "No rows in daily_metrics!"
assert row_count > 1000, f"Suspicious low row count: {row_count}"
return {"row_count": row_count}
run_dbt = BashOperator(
task_id="refresh_dbt_models",
bash_command="cd /dbt && dbt run --target production --models tag:daily",
)
@task()
def send_notification(quality_result: dict):
from airflow.providers.slack.hooks.slack import SlackHook
hook = SlackHook("slack_data_team")
hook.send(
channel="#pipeline-alerts",
text=f"✅ E-Commerce pipeline complete | {quality_result['row_count']} rows | {{ ds }}",
)
end = EmptyOperator(task_id="pipeline_end")
# Build the dependency graph
orders = extract_orders()
products = extract_products()
start >> [wait_for_file, orders, products]
transform = run_spark_transform(orders, products)
wait_for_file >> transform
quality = data_quality_check(transform)
quality >> run_dbt >> send_notification(quality) >> end
ecommerce_analytics()
💡 Key Takeaways from This Pipeline
- Uses
EmptyOperatorfor clear start/end markers - Deferrable sensor for long waits (clickstream data)
- Parallel extraction from multiple sources
- Data quality gate before loading to production
- Slack notification with row count for observability
- Documentation embedded in the DAG via
doc_md