DE Wikiguides / pipelines

Data Pipelines

A data pipeline is a series of processing steps that move data from one or more sources to a destination (data warehouse, data lake, analytics platform). Pipelines extract data, apply transformations, and load the results into storage for analysis.

Core Pipeline Architecture

Every data pipeline consists of a few fundamental components, regardless of the specific tools involved:

  • Sources - Databases (Postgres, MySQL), APIs, event streams (Kafka), file storage (S3), SaaS platforms
  • Extract - Read data from sources via connectors, change data capture (CDC), or batch exports
  • Transform - Clean, enrich, aggregate, and reshape data using dbt, Spark, SQL, or custom code
  • Load - Write transformed data to the target destination (warehouse, lake, database)
  • Orchestrate - Schedule, monitor, and manage pipeline execution with Airflow, Prefect, or Dagster

Pipeline Example: PostgreSQL to Snowflake

# Simple pipeline flow using Airflow + dbt

# Extract: pull orders from Postgres
pg_hook = PostgresHook(postgres_conn_id="source_db")
records = pg_hook.get_records("SELECT * FROM orders WHERE updated_at > {{ ds }}")

# Stage: write to S3 as Parquet
s3_hook = S3Hook(aws_conn_id="aws_default")
s3_hook.load_bytes(
    parquet_bytes(records),
    key=f"orders/{ds}.parquet",
    bucket="landing-zone"
)

# Load: COPY into Snowflake staging
snowflake_hook = SnowflakeHook(snowflake_conn_id="snowflake")
snowflake_hook.run(f"""
    COPY INTO stg_orders
    FROM @s3_stage/orders/{ds}.parquet
    FILE_FORMAT = (TYPE = PARQUET)
""")

# Transform: run dbt models for cleaning and aggregation
dbt_run(models=["stg_orders", "fct_orders", "dim_customers"])

Pipeline Design Patterns

1. Linear Pipeline

The simplest pattern: extract, transform, load execute sequentially. Suitable for simple ETL jobs with single-source inputs.

2. Fan-out / Fan-in

One extraction feeds multiple parallel transforms (fan-out), and their results merge into a single destination (fan-in). Common in multi-source aggregation pipelines.

3. Branching Pipelines

Different processing paths for different data types or SLAs. A pipeline might send real-time data to a streaming sink and batch data to a warehouse.

4. Event-driven Pipelines

Triggered by events rather than fixed schedules. A new file landing in S3 triggers an extract; a Kafka message triggers a transform. Tools: AWS Lambda, Kinesis, Airflow sensors.

Error Handling

Robust pipelines handle failures gracefully at every stage:

  • Retries - Transient failures (network, rate limits) should retry with exponential backoff
  • Dead-letter queues - Records that fail repeatedly are routed to a DLQ for manual inspection
  • Idempotency - Running the same pipeline twice produces identical results (upsert, drop-and-reload)
  • Checkpointing - Long-running pipelines save progress so they can resume from the last successful checkpoint
  • Notifications - Alert via Slack/email/PagerDuty on failure, success, or SLA breach

Monitoring and Observability

Key metrics to track for every pipeline:

  • Record count and volume at each stage
  • Latency (end-to-end and per-stage)
  • Error rate and failure distribution
  • Data quality checks (nulls, duplicates, freshness)
  • SLAs and service-level indicators

External Resources

Data Pipelines

A data pipeline is a series of processing steps that move data from one or more sources to a destination (data warehouse, data lake, analytics platform). Pipelines extract data, apply transformations, and load the results into storage for analysis.

Core Pipeline Architecture

Every data pipeline consists of a few fundamental components, regardless of the specific tools involved:

  • Sources - Databases (Postgres, MySQL), APIs, event streams (Kafka), file storage (S3), SaaS platforms
  • Extract - Read data from sources via connectors, change data capture (CDC), or batch exports
  • Transform - Clean, enrich, aggregate, and reshape data using dbt, Spark, SQL, or custom code
  • Load - Write transformed data to the target destination (warehouse, lake, database)
  • Orchestrate - Schedule, monitor, and manage pipeline execution with Airflow, Prefect, or Dagster

Pipeline Example: PostgreSQL to Snowflake

# Simple pipeline flow using Airflow + dbt

# Extract: pull orders from Postgres
pg_hook = PostgresHook(postgres_conn_id="source_db")
records = pg_hook.get_records("SELECT * FROM orders WHERE updated_at > {{ ds }}")

# Stage: write to S3 as Parquet
s3_hook = S3Hook(aws_conn_id="aws_default")
s3_hook.load_bytes(
    parquet_bytes(records),
    key=f"orders/{ds}.parquet",
    bucket="landing-zone"
)

# Load: COPY into Snowflake staging
snowflake_hook = SnowflakeHook(snowflake_conn_id="snowflake")
snowflake_hook.run(f"""
    COPY INTO stg_orders
    FROM @s3_stage/orders/{ds}.parquet
    FILE_FORMAT = (TYPE = PARQUET)
""")

# Transform: run dbt models for cleaning and aggregation
dbt_run(models=["stg_orders", "fct_orders", "dim_customers"])

Pipeline Design Patterns

1. Linear Pipeline

The simplest pattern: extract, transform, load execute sequentially. Suitable for simple ETL jobs with single-source inputs.

2. Fan-out / Fan-in

One extraction feeds multiple parallel transforms (fan-out), and their results merge into a single destination (fan-in). Common in multi-source aggregation pipelines.

3. Branching Pipelines

Different processing paths for different data types or SLAs. A pipeline might send real-time data to a streaming sink and batch data to a warehouse.

4. Event-driven Pipelines

Triggered by events rather than fixed schedules. A new file landing in S3 triggers an extract; a Kafka message triggers a transform. Tools: AWS Lambda, Kinesis, Airflow sensors.

Error Handling

Robust pipelines handle failures gracefully at every stage:

  • Retries - Transient failures (network, rate limits) should retry with exponential backoff
  • Dead-letter queues - Records that fail repeatedly are routed to a DLQ for manual inspection
  • Idempotency - Running the same pipeline twice produces identical results (upsert, drop-and-reload)
  • Checkpointing - Long-running pipelines save progress so they can resume from the last successful checkpoint
  • Notifications - Alert via Slack/email/PagerDuty on failure, success, or SLA breach

Monitoring and Observability

Key metrics to track for every pipeline:

  • Record count and volume at each stage
  • Latency (end-to-end and per-stage)
  • Error rate and failure distribution
  • Data quality checks (nulls, duplicates, freshness)
  • SLAs and service-level indicators

External Resources