home
diamond Go Premium
Data Engineering Path  ·  Airflow
Apache Airflow Logo

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 EmptyOperator for 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
lock

This content is reserved for Premium Members.

Upgrade to Premium

Entity Details

Create New Item

help

Submit Technical Query

Have a question or run into an issue? Describe it below, upload an optional screenshot, and our engineering team will answer it!

image Attach image (optional)

Submit Feedback

build Free Developer Utility Free Tool
gavel

Privacy & Legal Disclaimer

1. Client-Side Browser Processing

All utility tools on DeepEngineerHub (including Image to PDF, Text Formatters, JSON Converters, and Encryptors) execute 100% locally within your client browser using WebAssembly and JavaScript. No uploaded images, text, or documents are transmitted, collected, or stored on remote servers.

2. Limitation of Liability ("As-Is" Provision)

Tools and services are provided free of charge for convenience and educational purposes "as-is" without warranties of any kind. DeepEngineerHub shall not be held liable for any data loss, formatting inconsistencies, or indirect damages resulting from tool usage.

3. Open Source & Third-Party Software

Certain utilities utilize open-source client libraries (such as jsPDF, Mermaid.js, Pyodide) licensed under MIT, Apache, or BSD open licenses. All intellectual property remains with their respective copyright holders.