← Portfolio

Medallion lakehouse: Airflow, dbt and PySpark on Delta Lake

This pipeline builds a medallion lakehouse (bronze, silver and gold Delta tables on S3), orchestrated by Apache Airflow 3 on Kubernetes. PySpark jobs, started by Airflow as SparkApplications on the Spark Operator, load raw JSONL and CSV files from the landing zone into bronze; to stress-test the cluster, one job multiplies 3 million flights 330 times into 990 million rows. dbt (dbt-spark through the Spark Thrift Server) turns bronze into typed, deduplicated silver tables and builds the gold table, and every model runs its data tests in the same dbt build. A second DAG simulates a daily delivery: PySpark lands one day of 1 million flights in bronze and MERGEs it into silver, so a re-run updates instead of duplicating. GitLab CI converts the notebooks to Python and copies them and the dbt project to S3 after every merge, and Airflow pulls the code from there at run time.

Apache Airflow 3dbtPySparkDelta LakeMedallion architectureSpark OperatorKubernetesGitLab CI

Medallion architecture

Each layer is a set of Delta tables in the Hive Metastore. PySpark does the heavy loading, dbt owns the SQL transformations and tests.

Landing zone

raw files on S3
  • airline_fleet/*.jsonlfleet size per airline
  • flights_bigfile/*.csv3 million US flights

Bronze

PySpark · as delivered + load metadata
  • bronze.airline_fleet+ _ingested_at, _source_file
  • bronze.flights_bigfile_raw3 million rows
  • bronze.flights_bigfile_raw_exploded× 330 = 990 million rows (volume test)
  • bronze.flights_bigfile_raw_deltadaily delivery, 1 million rows

Silver

dbt + PySpark MERGE · typed, renamed, deduplicated
  • silver.airline_fleetdbt incremental MERGE on carrier_code
  • silver.flights_bigfile_silverdbt rebuild, dedup with GROUP BY + max_by; daily PySpark MERGE on the flight key

Gold

dbt · business-ready
  • gold.airline_dest_fleetflights per airline and destination, with fleet size
  • gold.arline_codesairline names (reference table)

Airflow DAG airline_fleet_pipeline: run replay

Graph view of the DAG. The replay plays back the real task durations of the 990-million-row run, compressed to about 25 seconds. A task starts as soon as all its upstream tasks have succeeded.

0:00:00
PySpark: SparkKubernetesOperator dbt build: KubernetesPodOperator running success

Airflow grid view: run history

All runs of both DAGs, as in Airflow's grid view: one column per run, the bar on top is the run duration. Tasks were added to the DAG over time, so earlier runs have empty cells. Hover for details.

success failed upstream failed task not in the DAG yet

dbt: silver and gold models

dbt runs in its own pod (KubernetesPodOperator) and executes SQL on Spark through the Thrift Server. Each Airflow task runs dbt build --select <model>: the model plus its tests.

-- silver/airline_fleet.sql: incremental MERGE
{{ config(
    materialized="incremental",
    incremental_strategy="merge",
    unique_key="carrier_code",
) }}

with typed as (
  select upper(trim(airline_code)) as carrier_code,
         trim(airline)              as carrier_name,
         cast(fleet_size as smallint) as aircraft_count,
         cast(_ingested_at as timestamp) as ingested_at
  from {{ source('bronze', 'airline_fleet') }}
)
-- MERGE fails on duplicate keys: keep the latest per airline
select * from (
  select *, row_number() over (
           partition by carrier_code order by ingested_at desc) as rn
  from typed) where rn = 1

Lesson learned: a MERGE of 990 million source rows into 990 million existing rows ran out of memory and disk on this small cluster. The big silver table is therefore rebuilt each run, and deduplicated with GROUP BY + max_by (hash aggregation) instead of row_number(), which has to sort every partition.

Airflow DAG flights_daily_delta: daily incremental load

Every run delivers the next day of 1 million flights. PySpark applies the same typing and deduplication as the dbt model, then MERGEs on the flight key: a new day is inserted, a re-run of the same day updates.

0:00:00

CI/CD: from git to the cluster

No code is baked into the Airflow image: tasks fetch their code from S3 when they start.

Merge to main→ GitLab CI→ notebooks .ipynb → .py+ dbt project→ S3 (Scaleway Object Storage)→ Spark driver / dbt pod pulls code at run time