How to Build Scalable ETL Pipelines with the Modern Data Stack: A Complete Guide
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, 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. For batch analytics, Apache Spark handles complex joins and aggregations, with example jobs like 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. 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. 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 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.
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 that reads from Kafka and writes to PostgreSQL. The pattern demonstrates exactly-once processing with SQL-based transformations.
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. The following pattern reads from PostgreSQL and writes to Delta Lake, enabling time-travel and ACID transactions.
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.
-- 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 to ensure data quality:
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.
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 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.pyandteam_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 to establish run-books and escalation procedures for data quality incidents.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →