- name
- datatalks-data-engineering-zoomcamp
- description
- Free 9-week data engineering course covering Docker, Terraform, Kestra, BigQuery, dbt, Spark, and Kafka with hands-on projects
- triggers
- ["help me with the data engineering zoomcamp","how do I set up the DE zoomcamp environment","show me how to complete zoomcamp homework","what are the data engineering zoomcamp modules","help with zoomcamp docker setup","configure terraform for zoomcamp GCP","run zoomcamp spark exercises","complete data engineering course project"]
# DataTalks Data Engineering Zoomcamp
> Skill by [ara.so](https://ara.so) — Data Skills collection.
## Overview
The Data Engineering Zoomcamp is a comprehensive 9-week free course covering production-ready data pipeline development. It includes hands-on modules on containerization (Docker), infrastructure as code (Terraform), workflow orchestration (Kestra), data warehousing (BigQuery), analytics engineering (dbt), data platforms (Bruin), batch processing (Spark), and streaming (Kafka).
The course operates in cohorts (next starts January 2026) but all materials are available for self-paced learning.
## Prerequisites
- Basic coding experience
- SQL familiarity
- Python knowledge (helpful but not required)
- Git installed
- Docker Desktop or Docker Engine
- Google Cloud Platform (GCP) account (free tier)
## Course Structure
### Module 1: Docker & Terraform
**Set up containerized PostgreSQL database:**
```bash
# Create network
docker network create pg-network
# Run PostgreSQL
docker run -d \
--name pg-database \
--network pg-network \
-e POSTGRES_USER=root \
-e POSTGRES_PASSWORD=root \
-e POSTGRES_DB=ny_taxi \
-v $(pwd)/ny_taxi_postgres_data:/var/lib/postgresql/data \
-p 5432:5432 \
postgres:13
# Run pgAdmin
docker run -d \
--name pgadmin \
--network pg-network \
-e PGADMIN_DEFAULT_EMAIL=admin@admin.com \
-e PGADMIN_DEFAULT_PASSWORD=root \
-p 8080:80 \
dpage/pgadmin4
```
**Docker Compose for entire stack:**
```yaml
# docker-compose.yaml
services:
pgdatabase:
image: postgres:13
environment:
- POSTGRES_USER=root
- POSTGRES_PASSWORD=root
- POSTGRES_DB=ny_taxi
volumes:
- ./ny_taxi_postgres_data:/var/lib/postgresql/data
ports:
- "5432:5432"
pgadmin:
image: dpage/pgadmin4
environment:
- PGADMIN_DEFAULT_EMAIL=admin@admin.com
- PGADMIN_DEFAULT_PASSWORD=root
ports:
- "8080:80"
```
```bash
# Start services
docker-compose up -d
# Stop services
docker-compose down
```
**Terraform GCP setup:**
```hcl
# main.tf
terraform {
required_version = ">= 1.0"
backend "local" {}
required_providers {
google = {
source = "hashicorp/google"
}
}
}
provider "google" {
project = var.project
region = var.region
}
# Data Lake Bucket
resource "google_storage_bucket" "data-lake-bucket" {
name = "${local.data_lake_bucket}_${var.project}"
location = var.region
storage_class = var.storage_class
uniform_bucket_level_access = true
versioning {
enabled = true
}
lifecycle_rule {
action {
type = "Delete"
}
condition {
age = 30
}
}
force_destroy = true
}
# BigQuery Dataset
resource "google_bigquery_dataset" "dataset" {
dataset_id = var.BQ_DATASET
project = var.project
location = var.region
}
```
```hcl
# variables.tf
locals {
data_lake_bucket = "dtc_data_lake"
}
variable "project" {
description = "Your GCP Project ID"
}
variable "region" {
description = "Region for GCP resources"
default = "europe-west6"
type = string
}
variable "storage_class" {
description = "Storage class type for your bucket"
default = "STANDARD"
}
variable "BQ_DATASET" {
description = "BigQuery Dataset"
type = string
default = "trips_data_all"
}
```
```bash
# Initialize Terraform
terraform init
# Plan infrastructure
terraform plan
# Apply infrastructure
terraform apply
# Destroy infrastructure
terraform destroy
```
### Module 2: Workflow Orchestration (Kestra)
**Example Kestra workflow for data ingestion:**
```yaml
# flows/ingest_data.yaml
id: ingest_ny_taxi_data
namespace: zoomcamp
tasks:
- id: download_data
type: io.kestra.core.tasks.scripts.Bash
commands:
- wget https://github.com/DataTalksClub/nyc-tlc-data/releases/download/yellow/yellow_tripdata_2021-01.csv.gz
- gunzip yellow_tripdata_2021-01.csv.gz
- id: python_ingest
type: io.kestra.plugin.scripts.python.Script
docker:
image: python:3.9
script: |
import pandas as pd
from sqlalchemy import create_engine
import os
df = pd.read_csv('yellow_tripdata_2021-01.csv', nrows=100000)
engine = create_engine(os.getenv('POSTGRES_CONNECTION'))
df.to_sql('yellow_taxi_data', engine, if_exists='replace', chunksize=10000)
print(f"Inserted {len(df)} rows")
- id: log_completion
type: io.kestra.core.tasks.log.Log
message: "Data ingestion completed successfully"
```
**Python data ingestion script:**
```python
# ingest_data.py
import pandas as pd
from sqlalchemy import create_engine
import argparse
from time import time
def main(params):
user = params.user
password = params.password
host = params.host
port = params.port
db = params.db
table_name = params.table_name
url = params.url
# Download CSV
csv_name = 'output.csv'
os.system(f"wget {url} -O {csv_name}")
# Create engine
engine = create_engine(f'postgresql://{user}:{password}@{host}:{port}/{db}')
# Read CSV in chunks
df_iter = pd.read_csv(csv_name, iterator=True, chunksize=100000)
df = next(df_iter)
# Convert datetime columns
df.tpep_pickup_datetime = pd.to_datetime(df.tpep_pickup_datetime)
df.tpep_dropoff_datetime = pd.to_datetime(df.tpep_dropoff_datetime)
# Create table
df.head(n=0).to_sql(name=table_name, con=engine, if_exists='replace')
# Insert first chunk
df.to_sql(name=table_name, con=engine, if_exists='append')
# Insert remaining chunks
while True:
try:
t_start = time()
df = next(df_iter)
df.tpep_pickup_datetime = pd.to_datetime(df.tpep_pickup_datetime)
df.tpep_dropoff_datetime = pd.to_datetime(df.tpep_dropoff_datetime)
df.to_sql(name=table_name, con=engine, if_exists='append')
t_end = time()
print(f'Inserted another chunk, took %.3f seconds' % (t_end - t_start))
except StopIteration:
print("Finished ingesting data")
break
if __name__ == '__main__':
parser = argparse.ArgumentParser(description='Ingest CSV data to Postgres')
parser.add_argument('--user', required=True, help='user name for postgres')
parser.add_argument('--password', required=True, help='password for postgres')
parser.add_argument('--host', required=True, help='host for postgres')
parser.add_argument('--port', required=True, help='port for postgres')
parser.add_argument('--db', required=True, help='database name for postgres')
parser.add_argument('--table_name', required=True, help='name of the table')
parser.add_argument('--url', required=True, help='url of the csv file')
args = parser.parse_args()
main(args)
```
```bash
# Run ingestion script
python ingest_data.py \
--user=root \
--password=root \
--host=localhost \
--port=5432 \
--db=ny_taxi \
--table_name=yellow_taxi_trips \
--url=https://github.com/DataTalksClub/nyc-tlc-data/releases/download/yellow/yellow_tripdata_2021-01.csv.gz
```
### Module 3: Data Warehouse (BigQuery)
**Create partitioned and clustered table:**
```sql
-- Create external table
CREATE OR REPLACE EXTERNAL TABLE `trips_data_all.external_yellow_tripdata`
OPTIONS (
format = 'CSV',
uris = ['gs://nyc-tl-data/trip data/yellow_tripdata_2019-*.csv',
'gs://nyc-tl-data/trip data/yellow_tripdata_2020-*.csv']
);
-- Create partitioned table
CREATE OR REPLACE TABLE `trips_data_all.yellow_tripdata_partitioned`
PARTITION BY
DATE(tpep_pickup_datetime) AS
SELECT * FROM `trips_data_all.external_yellow_tripdata`;
-- Create partitioned and clustered table
CREATE OR REPLACE TABLE `trips_data_all.yellow_tripdata_partitioned_clustered`
PARTITION BY DATE(tpep_pickup_datetime)
CLUSTER BY VendorID AS
SELECT * FROM `trips_data_all.external_yellow_tripdata`;
-- Query comparison
SELECT DISTINCT(VendorID)
FROM `trips_data_all.yellow_tripdata_partitioned`
WHERE DATE(tpep_pickup_datetime) BETWEEN '2020-06-01' AND '2020-06-30';
-- This query will process less data vs non-partitioned
```
**Load data from GCS to BigQuery:**
```python
# load_to_bq.py
from google.cloud import bigquery
import os
# Set credentials
os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = 'path/to/credentials.json'
# Initialize client
client = bigquery.Client()
# Define table
table_id = 'your-project.trips_data_all.yellow_tripdata'
# Configure load job
job_config = bigquery.LoadJobConfig(
source_format=bigquery.SourceFormat.PARQUET,
write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE,
)
uri = 'gs://your-bucket/yellow_tripdata_2021-01.parquet'
# Load data
load_job = client.load_table_from_uri(
uri, table_id, job_config=job_config
)
load_job.result() # Wait for job to complete
print(f"Loaded {load_job.output_rows} rows to {table_id}")
```
### Module 4: Analytics Engineering (dbt)
**Project structure:**
```
dbt_project/
├── dbt_project.yml
├── profiles.yml
├── models/
│ ├── staging/
│ │ ├── stg_yellow_tripdata.sql
│ │ └── schema.yml
│ └── core/
│ ├── fact_trips.sql
│ └── dim_zones.sql
└── macros/
└── get_payment_type_description.sql
```
**dbt_project.yml:**
```yaml
name: 'taxi_rides_ny'
version: '1.0.0'
config-version: 2
profile: 'default'
model-paths: ["models"]
analysis-paths: ["analyses"]
test-paths: ["tests"]
seed-paths: ["seeds"]
macro-paths: ["macros"]
snapshot-paths: ["snapshots"]
target-path: "target"
clean-targets:
- "target"
- "dbt_packages"
models:
taxi_rides_ny:
staging:
+materialized: view
core:
+materialized: table
```
**profiles.yml:**
```yaml
default:
outputs:
dev:
type: bigquery
method: service-account
project: "{{ env_var('GCP_PROJECT_ID') }}"
dataset: dbt_dev
threads: 4
keyfile: "{{ env_var('GOOGLE_APPLICATION_CREDENTIALS') }}"
location: EU
prod:
type: bigquery
method: service-account
project: "{{ env_var('GCP_PROJECT_ID') }}"
dataset: production
threads: 4
keyfile: "{{ env_var('GOOGLE_APPLICATION_CREDENTIALS') }}"
location: EU
target: dev
```
**Staging model (models/staging/stg_yellow_tripdata.sql):**
```sql
{{ config(materialized='view') }}
with tripdata as
(
select *,
row_number() over(partition by vendorid, tpep_pickup_datetime) as rn
from {{ source('staging','yellow_tripdata') }}
where vendorid is not null
)
select
-- identifiers
{{ dbt_utils.generate_surrogate_key(['vendorid', 'tpep_pickup_datetime']) }} as tripid,
cast(vendorid as integer) as vendorid,
cast(ratecodeid as integer) as ratecodeid,
cast(pulocationid as integer) as pickup_locationid,
cast(dolocationid as integer) as dropoff_locationid,
-- timestamps
cast(tpep_pickup_datetime as timestamp) as pickup_datetime,
cast(tpep_dropoff_datetime as timestamp) as dropoff_datetime,
-- trip info
store_and_fwd_flag,
cast(passenger_count as integer) as passenger_count,
cast(trip_distance as numeric) as trip_distance,
-- payment info
cast(fare_amount as numeric) as fare_amount,
cast(extra as numeric) as extra,
cast(mta_tax as numeric) as mta_tax,
Ver no GitHub