# How to Build Scalable ETL Pipelines with the Modern Data Stack: A Complete Guide

> Build scalable ETL pipelines using Kafka, Flink/Spark, Delta Lake, dbt, and Airflow. Master modern data stack tools for efficient data processing and robust pipelines. Get the complete guide.

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

---

**Build scalable ETL pipelines by combining Kafka for ingestion, Apache Flink or Spark for processing, Delta Lake for storage, dbt for modeling, and Airflow for orchestration.**

The **DataExpert-io/data-engineer-handbook** repository provides hands-on examples for constructing production-grade data pipelines using best-in-class tools from the modern data stack. By decoupling ingestion, transformation, and storage into specialized layers, you can create fault-tolerant systems that handle increasing data volumes without compromising performance. This guide walks through the architecture and implementation patterns found in the handbook, showing you exactly how to assemble each component into a cohesive, scalable ETL pipeline.

## Architectural Layers of the Modern Data Stack

Scalable pipelines rely on five distinct layers, each optimized for a specific function. The handbook organizes its examples around this architecture, from real-time ingestion to observability.

### Ingestion (Extract)

Modern pipelines start with **Kafka** for high-throughput event streaming or managed tools like **Airbyte** and **Fivetran** for SaaS connectors. According to the [`intermediate-bootcamp/materials/4-apache-flink-training/README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/README.md), the repository uses a Confluent Cloud Kafka cluster to ingest web-event data, with credentials configured in `flink-env.env`. Kafka partitions enable parallel consumption, allowing the extraction layer to scale horizontally as event volume grows.

### Processing (Transform)

The handbook demonstrates two primary processing paradigms. For real-time workloads, **Apache Flink** provides exactly-once processing guarantees, as shown in [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py). For batch analytics, **Apache Spark** handles complex joins and aggregations, with example jobs like [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py) located in `intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/`. **dbt** complements these engines by enabling version-controlled, testable SQL transformations.

### Storage (Load)

The repository recommends **Delta Lake** and **Apache Iceberg** for data lake architectures, listed under the "Data Lake / Cloud" section of the main [`README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/README.md). For analytical workloads, the Flink tutorial writes transformed events into a PostgreSQL `processed_events` table. Columnar formats like Parquet combined with transaction logs enable incremental reads and ACID guarantees essential for scalable pipelines.

### Orchestration and Scheduling

Production pipelines require robust DAG management. The handbook lists **Apache Airflow**, **Dagster**, **Prefect**, and **Kestra** under the "Orchestration" section of [`README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/README.md). These tools support dynamic task generation, containerized execution, and native integration with Spark and Flink clusters.

### Observability and Quality

The [`intermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md) details run-books and on-call schedules for production support. For data quality, the repository recommends **Great Expectations**, **dbt tests**, and **OpenLineage** to catch schema drift early and provide impact analysis for downstream model changes.

## Implementing the Extract Layer with Kafka

To simulate real-time ingestion, the handbook provides a Kafka producer pattern that streams JSON events to a Confluent Cloud topic. This matches the source used in the Flink streaming exercises.

```python
from confluent_kafka import Producer
import json, time, uuid

p = Producer({
    'bootstrap.servers': 'pkc-rgm37.us-west-2.aws.confluent.cloud:9092',
    'security.protocol': 'SASL_SSL',
    'sasl.mechanisms': 'PLAIN',
    'sasl.username': '<API_KEY>',
    'sasl.password': '<API_SECRET>'
})

def delivery_report(err, msg):
    if err is not None:
        print(f'Delivery failed: {err}')
    else:
        print(f'Message delivered to {msg.topic()} [{msg.partition()}]')

while True:
    event = {
        "event_id": str(uuid.uuid4()),
        "event_type": "click",
        "timestamp": int(time.time()*1000),
        "user_id": f"user_{int(time.time())%1000}"
    }
    p.produce('bootcamp-events-prod', json.dumps(event).encode('utf-8'), callback=delivery_report)
    p.poll(0)
    time.sleep(0.5)

```

This producer writes to the `bootcamp-events-prod` topic referenced in the Flink training materials, establishing the data source for downstream transformation.

## Real-Time Transformation with Apache Flink

For stream processing, the handbook implements a PyFlink job in [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py) that reads from Kafka and writes to PostgreSQL. The pattern demonstrates exactly-once processing with SQL-based transformations.

```python
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes, EnvironmentSettings, Schema
from pyflink.table.descriptors import Kafka, Json, Jdbc

env = StreamExecutionEnvironment.get_execution_environment()
settings = EnvironmentSettings.new_instance().in_streaming_mode().use_blink_planner().build()
t_env = StreamTableEnvironment.create(env, environment_settings=settings)

# Define Kafka source

t_env.connect(
    Kafka()
    .version("universal")
    .topic("bootcamp-events-prod")
    .property("bootstrap.servers", "pkc-rgm37.us-west-2.aws.confluent.cloud:9092")
    .property("sasl.mechanism", "PLAIN")
    .property("security.protocol", "SASL_SSL")
    .property("sasl.jaas.config",
              "org.apache.kafka.common.security.plain.PlainLoginModule required username='<API_KEY>' password='<API_SECRET>';")
    .start_from_earliest()
).with_format(
    Json().fail_on_missing_field(False)
).with_schema(
    Schema()
    .field("event_id", DataTypes.STRING())
    .field("event_type", DataTypes.STRING())
    .field("timestamp", DataTypes.BIGINT())
    .field("user_id", DataTypes.STRING())
).create_temporary_table("kafka_source")

# Define PostgreSQL sink

t_env.connect(
    Jdbc()
    .url("jdbc:postgresql://host.docker.internal:5432/postgres")
    .table("processed_events")
    .driver("org.postgresql.Driver")
    .username("postgres")
    .password("postgres")
).with_format(
    Json()
).with_schema(
    Schema()
    .field("event_id", DataTypes.STRING())
    .field("event_type", DataTypes.STRING())
    .field("timestamp", DataTypes.BIGINT())
    .field("user_id", DataTypes.STRING())
).create_temporary_table("postgres_sink")

# Simple transformation: filter only click events

t_env.sql_query("""
    SELECT *
    FROM kafka_source
    WHERE event_type = 'click'
""").insert_into("postgres_sink")

t_env.execute("flink_kafka_to_postgres")

```

This job filters click events in real-time before loading them into the `processed_events` table, demonstrating low-latency transformation with exactly-once guarantees.

## Batch Processing with Apache Spark

For larger-scale aggregations, the handbook includes PySpark examples such as [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py). The following pattern reads from PostgreSQL and writes to Delta Lake, enabling time-travel and ACID transactions.

```python
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("PostgresToDelta") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# Load raw events

df = spark.read \
    .format("jdbc") \
    .option("url", "jdbc:postgresql://host.docker.internal:5432/postgres") \
    .option("dbtable", "processed_events") \
    .option("user", "postgres") \
    .option("password", "postgres") \
    .load()

# Simple enrichment

df_enriched = df.withColumn("event_date", df.timestamp.cast("timestamp").cast("date"))

# Write to Delta Lake (partitioned by date)

df_enriched.write \
    .format("delta") \
    .partitionBy("event_date") \
    .mode("append") \
    .save("/mnt/delta/events")

```

Partitioning by `event_date` optimizes query performance for time-series analytics while Delta Lake's transaction log ensures data consistency during concurrent writes.

## Data Modeling with dbt

The handbook recommends **dbt** for declarative transformations. While the repository does not contain a full dbt project, the README lists it under "Data Quality" tools. The following example shows a staging model with built-in testing.

```sql
-- models/stg_events.sql
with source as (
    select *
    from {{ ref('raw_events') }}
),

clean as (
    select
        event_id,
        user_id,
        cast(timestamp as timestamp) as event_ts,
        event_type
    from source
    where event_type = 'click'
)

select *
from clean

```

Configure tests in [`dbt_project.yml`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/dbt_project.yml) to ensure data quality:

```yaml
models:
  my_project:
    stg_events:
      materialized: view
      tests:
        - unique:
            column_name: event_id
        - not_null:
            column_name: user_id

```

This separates transformation logic from execution, enabling version control and automated testing of business logic.

## Orchestrating Pipelines with Apache Airflow

The handbook lists Airflow among essential orchestration tools. The following DAG coordinates the Flink, Spark, and dbt components into a unified workflow.

```python
from airflow import DAG
from airflow.providers.docker.operators.docker import DockerOperator
from airflow.utils.dates import days_ago

default_args = {
    "owner": "data-engineer",
    "retries": 2,
}

with DAG(
    dag_id="etl_modern_stack",
    default_args=default_args,
    schedule_interval="@daily",
    start_date=days_ago(1),
    catchup=False,
) as dag:

    # Run Flink job (containerized)

    run_flink = DockerOperator(
        task_id="run_flink",
        image="my-flink-image:latest",
        command="python /opt/src/job/start_job.py",
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
    )

    # Spark batch job

    run_spark = DockerOperator(
        task_id="run_spark",
        image="my-spark-image:latest",
        command="spark-submit /opt/src/jobs/weekly_aggregation.py",
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
    )

    # dbt transformation

    run_dbt = DockerOperator(
        task_id="run_dbt",
        image="ghcr.io/dbt-labs/dbt:latest",
        command="dbt run --profiles-dir /opt/profiles",
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
    )

    run_flink >> run_spark >> run_dbt

```

This containerized approach ensures consistent environments across development and production while handling dependency management and retry logic.

## Ensuring Reliability and Observability

Production maintenance requires structured operational practices. The [`intermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md) emphasizes run-books and on-call schedules for incident response. Combine these operational practices with **Great Expectations** for automated data validation and **OpenLineage** for dependency tracking to create a fully observable pipeline that scales reliably.

## Summary

- **Decouple layers** using Kafka for ingestion, Flink/Spark for processing, and Delta Lake for storage to enable independent scaling of each component.
- **Implement exactly-once processing** with Apache Flink for real-time streams and Apache Spark for batch workloads, as demonstrated in [`start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/start_job.py) and [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py).
- **Use dbt** to version-control SQL transformations and enforce data quality through automated testing.
- **Orchestrate with Airflow** to manage dependencies between containerized tasks and handle failure recovery.
- **Maintain operational readiness** by following the run-book and on-call guidelines in the pipeline maintenance module.

## Frequently Asked Questions

### What makes the modern data stack different from traditional ETL?

Traditional ETL typically uses monolithic, proprietary tools that couple extraction, transformation, and loading into a single platform. The modern data stack decouples these functions into best-in-class, open-source components—such as Kafka for ingestion, Spark for processing, and dbt for modeling—allowing each layer to scale independently according to specific workload requirements.

### How do I choose between Apache Flink and Apache Spark for my pipeline?

Choose **Apache Flink** when you need low-latency stream processing with exactly-once guarantees, as shown in the handbook's Kafka-to-PostgreSQL example. Use **Apache Spark** for complex batch aggregations, large-scale joins, or when you need mature ecosystem support for machine learning libraries. Many organizations use both: Flink for real-time ingestion and Spark for downstream batch analytics.

### What is the role of Delta Lake in scalable ETL pipelines?

Delta Lake adds an ACID transaction log and time-travel capabilities to cloud storage systems like S3 or ADLS. This enables safe concurrent writes, schema enforcement, and incremental processing—essential features when multiple pipeline stages read and write the same datasets simultaneously. The handbook lists Delta Lake as a core component of the modern data lake architecture.

### How do I monitor data quality in a production pipeline?

Implement **Great Expectations** or **dbt tests** to validate schema, freshness, and distribution assumptions at runtime. The handbook recommends pairing these tools with **OpenLineage** for dependency tracking. Additionally, follow the maintenance practices in [`6-data-pipeline-maintenance/README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/6-data-pipeline-maintenance/README.md) to establish run-books and escalation procedures for data quality incidents.