# How to Build End-to-End Data Pipelines Using Apache Spark and Delta Lake: A Complete Guide

> Learn to build end-to-end data pipelines with Apache Spark and Delta Lake. Master streaming ingestion, ACID storage, and optimized transformations for robust data solutions.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: how-to-guide
- Published: 2026-08-12

---

**You can build end-to-end data pipelines using Apache Spark and Delta Lake by combining Structured Streaming for ingestion, Delta tables for ACID-compliant storage, and optimized Spark SQL transformations with broadcast and bucketed joins, orchestrated via tools like Dagster.**

The DataExpert-io/data-engineer-handbook repository provides a battle-tested blueprint for constructing production-grade data pipelines that leverage Apache Spark's distributed processing alongside Delta Lake's transactional guarantees. Whether you are ingesting raw web logs from S3 or implementing slowly changing dimensions, understanding how to build end-to-end data pipelines using Apache Spark and Delta Lake is essential for modern lakehouse architectures.

## Architecture Overview: The Five-Layer Pipeline

Robust pipelines follow a layered architecture that separates concerns from ingestion to consumption. According to the repository's practical project documentation in [`projects.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/projects.md), a complete stack flows from **Web-scraping → S3 → Spark & Delta Lake → Jupyter → Druid → Superset → Dagster**.

The architectural layers break down as follows:

- **Ingestion Layer**: Use Spark Structured Streaming to pull raw data from sources like S3, Kafka, or files into a landing Delta table.
- **Storage Layer**: Delta Lake serves as the unified data lake, housing raw, curated, and presentation layers as versioned tables with ACID guarantees.
- **Transformation Layer**: Execute joins, aggregations, and enrichment using Spark SQL or PySpark, optimizing with broadcast and bucketed join strategies.
- **Orchestration Layer**: Tools like Dagster schedule jobs, enforce dependencies, and monitor pipeline health.
- **Consumption Layer**: Queryable Delta tables feed downstream analytics, machine learning models, and BI tools.

## Core Spark Patterns for Delta Lake Performance

Optimizing Spark jobs requires explicit control over join strategies and write semantics. The handbook's Spark Fundamentals homework and production job files demonstrate specific patterns for Delta Lake integration.

### Disable Automatic Broadcast Joins for Predictability

Prevent Spark from unexpectedly broadcasting large tables by disabling the automatic threshold. This configuration is explicitly recommended in [`intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md) at lines 15-17:

```python
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

```

Setting this to `-1` forces you to declare broadcast intentions explicitly, avoiding out-of-memory errors on large datasets.

### Explicit Broadcast Joins for Dimension Tables

Force small dimension tables to broadcast, eliminating expensive shuffle operations. The [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py) file implements this pattern using the `broadcast` hint:

```python
from pyspark.sql.functions import broadcast

broadcast_df = broadcast(spark.table("medals"))

```

Broadcasting ensures the entire small table is sent to each executor, keeping the join local and fast.

### Bucketed Joins for Large Fact Tables

Pre-bucket large fact tables on frequently joined keys to speed up joins and aggregations. This pattern appears in [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py) for optimizing vertex table creation:

```python
df.repartition(16, "match_id").write.format("delta").saveAsTable("fact_bucketed")

```

Aligning the number of buckets with common join keys (like `match_id`) minimizes data movement across the cluster.

### Atomic Writes and ACID Guarantees

Write curated results using Delta's transactional guarantees. The repository demonstrates atomic overwrites using `insertInto`:

```python
output_df.write.format("delta").mode("overwrite").insertInto("team_vertex")

```

This ensures downstream queries see a consistent snapshot, never partial data.

## Implementing the Pipeline: From Raw to Curated

Below is a minimal, runnable pipeline based on the handbook's patterns. This example configures the Spark session with Delta extensions, ingests JSON logs from S3, applies optimized transformations, and writes to a curated Delta table:

```python

# spark_delta_pipeline.py

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, broadcast

# 1️⃣ Spark session with Delta Lake support

spark = (
    SparkSession.builder
    .appName("DeltaPipeline")
    .master("local[*]")
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
    .getOrCreate()
)

# 2️⃣ Ingest raw data from S3 landing zone

raw_df = spark.read.json("s3a://my-bucket/raw-logs/*.json")

# 3️⃣ Write raw layer as Delta (landing zone)

raw_df.write.format("delta").mode("overwrite").saveAsTable("raw_logs")

# 4️⃣ Transformations with optimized joins

#   a) Broadcast small lookup table

countries = spark.read.format("delta").load("s3a://my-bucket/lookups/countries")
countries = broadcast(countries)

#   b) Join and aggregate with bucketing

events = spark.read.format("delta").table("raw_logs")
joined = events.join(countries, "country_id")
aggregated = (
    joined.groupBy("event_type")
    .agg({"value": "sum"})
    .repartition(16, "event_type")  # bucketed partitioning

)

# 5️⃣ Write curated layer (Delta)

aggregated.write.format("delta").mode("overwrite").saveAsTable("curated_events")

```

This implementation demonstrates the **Delta IO extensions** required for table support, **explicit broadcast** for dimension tables, **bucketed repartition** for aggregation efficiency, and **atomic writes** for data consistency.

## Production-Ready Patterns from the Repository

The DataExpert-io/data-engineer-handbook contains specific job implementations that extend the basic patterns above into production workflows:

- **[`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py)**: Located at [`intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/team_vertex_job.py), this script demonstrates combined broadcast and bucketed join strategies for graph analytics vertex generation.
- **[`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py)**: Found at [`intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/players_scd_job.py), this implements Slowly Changing Dimension (SCD) Type 2 logic on Delta tables, handling historical record tracking with merge operations.
- **[`monthly_user_site_hits_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/monthly_user_site_hits_job.py)**: Available at [`intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/monthly_user_site_hits_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/monthly_user_site_hits_job.py), this shows time-window aggregations and Delta writes optimized for web analytics data.
- **[`projects.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/projects.md)**: The practical project guide at the repository root ([`projects.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/projects.md)) illustrates the full-stack integration connecting web scraping, S3, Spark, Delta Lake, Jupyter, Druid, Superset, and Dagster into a cohesive pipeline.

## Summary

Building reliable data pipelines requires combining the right architectural patterns with Delta Lake's ACID capabilities:

- **Configure Spark** with Delta extensions (`io.delta.sql.DeltaSparkSessionExtension`) to enable table support.
- **Optimize joins** by disabling auto-broadcast, explicitly broadcasting small dimensions, and bucketing large fact tables on join keys.
- **Implement medallion architecture** (raw → curated → presentation) using Delta tables for each layer to leverage time-travel and schema evolution.
- **Write atomically** using `mode("overwrite")` with `insertInto` or `saveAsTable` to ensure consistent downstream reads.
- **Orchestrate** with Dagster or similar tools, referencing the complete pipeline examples in the handbook's intermediate bootcamp materials.

## Frequently Asked Questions

### How do I enable Delta Lake support in a Spark session?

You must configure two essential Spark settings before interacting with Delta tables. Set `spark.sql.extensions` to `io.delta.sql.DeltaSparkSessionExtension` and `spark.sql.catalog.spark_catalog` to `org.apache.spark.sql.delta.catalog.DeltaCatalog` in your SparkSession builder, as demonstrated in the pipeline implementation above. These configurations register the Delta catalog and SQL extensions required for table management.

### What is the difference between broadcast and bucketed joins in Spark?

**Broadcast joins** copy the entire small table to every executor, ideal for dimension tables under 10MB, eliminating shuffle entirely. **Bucketed joins** pre-partition both tables on the join key into the same number of buckets, minimizing shuffle for large fact-to-fact joins. According to the handbook's homework in [`intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md), you should disable automatic broadcasting and explicitly choose the strategy based on table size and join frequency.

### How does Delta Lake handle schema evolution in production pipelines?

Delta Lake supports **schema enforcement** by default, rejecting incompatible writes to prevent data corruption. You can enable **schema evolution** using the `.option("mergeSchema", "true")` flag when writing, or use the `ALTER TABLE` command to add columns explicitly. This flexibility, combined with **time-travel** capabilities (querying previous table versions via `VERSION AS OF` or `TIMESTAMP AS OF`), allows pipelines to adapt to changing source schemas without breaking downstream consumers.

### Where can I find a complete example of an end-to-end pipeline using these technologies?

The [`projects.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/projects.md) file in the DataExpert-io/data-engineer-handbook repository documents a comprehensive pipeline spanning web scraping, S3 ingestion, Spark processing with Delta Lake storage, Jupyter analysis, Druid indexing, Superset visualization, and Dagster orchestration. Additionally, the `intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/` directory contains runnable Python scripts ([`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py), [`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py), [`monthly_user_site_hits_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/monthly_user_site_hits_job.py)) that demonstrate specific production patterns like SCD handling and bucketing strategies.