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