This project implements a containerized batch data engineering pipeline using Apache Airflow and Postgres.
It demonstrates append-only raw ingestion, batch semantics, and layered warehouse modeling (raw → staging).
Each Airflow DAG run simulates the ingestion of a new batch of source data and persists it immutably in the warehouse.
- Postgres 15
airflowdatabase → Airflow metadata (DAGs, task state, auth)de_dwdatabase → Analytics warehouse- Schemas:
raw,stg
- Schemas:
- Apache Airflow 2.8.3
- Executor:
LocalExecutor - Services:
- Webserver
- Scheduler
- Init container (DB bootstrap)
- Executor:
All services are orchestrated via Docker Compose.
Raw tables are append-only and store data exactly as received, including duplicates across batches.
Example:
raw.orders- Allows duplicate
order_idvalues - Each batch is identified by a shared
updated_attimestamp
- Allows duplicate
Raw data is never mutated or deduplicated.
The staging layer applies business logic to raw data.
Example:
stg.orders_latest- One row per
order_id - Keeps the latest record by
updated_at
- One row per
This separation ensures:
- Full historical traceability in raw
- Clean, analytics-ready data in staging
Each DAG run performs the following steps:
-
Generate batch data
- Create a synthetic batch of 100 orders
- Assign a single batch timestamp
- Write to
data/raw/orders.csv(overwritten each run)
-
Load raw data
- Read CSV
- Append rows to
raw.orders - No deduplication
-
Build staging model
- Deduplicate raw data
- Materialize
stg.orders_latest
Running the DAG multiple times produces multiple immutable batches in the raw layer.
After 3 DAG runs:
-
raw.orders- 300 total rows
- 3 distinct batch timestamps
- 100 distinct
order_idvalues
-
stg.orders_latest- 100 rows
- One per
order_id - All from the latest batch
- Docker
- Docker Compose
docker compose up -d