Roman Klimenko
EN

Syv måder at bygge en medallion-arkitektur i Databricks

Databricks giver dig flere måder at bygge en Bronze-, Silver- og Gold-pipeline på. Hver tilgang har sine egne afvejninger, men dokumentationen sammenligner dem sjældent direkte.

Jeg byggede den samme pipeline på syv forskellige måder og sammenlignede resultaterne. Alle versioner bruger de samme CSV-data og det samme skema og producerer det samme Kimball star schema. Projektet er deployet som et Databricks Asset Bundle, og koden ligger i dette repo.

De næste afsnit går fra den mest manuelle til den mest deklarative tilgang. Til sidst opsummerer jeg, hvornår jeg ville vælge hver af dem.

Testcasen

Datasættet er bevidst lille. Det indeholder kunder, produkter, ordrer og ordrelinjer fordelt på to batches. Batch 1 er den første indlæsning. I batch 2 ændrer to kunder deres oplysninger, en ny kunde bliver tilføjet, og en ny ordre dukker op. Det sætter SCD2-logikken på prøve.

Gold-laget er en klassisk dimensionel Kimball-model med dim_customer inklusive SCD2-historik, dim_product, dim_date og fact_order_line. Når begge batches er behandlet, skal hver tilgang producere præcis 8 kunderækker, nemlig 6 aktuelle inklusive den nye kunde og 2 historiske, samt 5 produkter, 91 datoer fra januar til marts 2024 og 11 fact-rækker.

Hvis tallene ikke stemmer, er der noget galt med implementeringen. Den begrænsning holder sammenligningen ærlig.

1. Python notebooks: den manuelle baseline

Kode: bronze.py, silver.py, gold.py

Denne tilgang bruger tre PySpark-notebooks, én til hvert lag. Bronze læser CSV-filer og tilføjer metadata. Silver deduplikerer rækkerne med row_number over et window, så den seneste batch vinder, og caster derefter typer og standardiserer tekst. Gold bygger den dimensionelle model.

window = (Window
    .partitionBy("customer_id")
    .orderBy(F.col("_batch_id").desc()))

silver_customers = (bronze_customers
    .withColumn("rn", F.row_number().over(window))
    .where("rn = 1"))

SCD2-logikken i Gold er den sværeste del. Du skal selv sammenligne batch-snapshots, bygge historiske rækker og vedligeholde valid_from, valid_to og is_current. Det virker, men koden er lang og nem at lave fejl i.

Du får fuld kontrol, men skal også skrive og vedligeholde logik, som platformen ellers kunne håndtere.

2. SQL notebooks med COPY INTO

Kode: bronze.sql, silver.sql, gold.sql

Denne version beholder samme arkitektur, men bruger SQL hele vejen. Bronze bruger COPY INTO i stedet for spark.read. Silver bruger CREATE OR REPLACE TABLE AS SELECT med CTE-baseret deduplikering. Gold bygger SCD2-historik med LEFT ANTI JOIN og change detection. Logikken er stadig manuel, men kan være lettere at læse for et team, der primært arbejder i SQL.

COPY INTO bronze_orders
FROM '/Volumes/.../orders'
    FILEFORMAT = CSV
FORMAT_OPTIONS ('header' = 'true');

For SQL-tunge teams føles denne tilgang naturlig og kan gøre code reviews lettere. Begrænsningerne er de samme: orkestrering sker opgave for opgave, inkrementel behandling er manuel, og SCD2 kræver omhyggelig merge-logik.

3. Materialized Views + Streaming Tables

Kode: setup.sql, scd2_merge.sql

Materialized views og streaming tables er mere deklarative. Du beskriver, hvad hver tabel skal indeholde, og Databricks finder ud af, hvordan den skal bygges.

CREATE OR REFRESH STREAMING TABLE bronze_orders AS
SELECT *
FROM STREAM read_files('/Volumes/.../orders', format => 'csv');

CREATE OR REPLACE MATERIALIZED VIEW silver_orders AS
SELECT ...
FROM bronze_orders;

Bronze-tabellerne bliver streaming tables baseret på Auto Loader, som holder styr på de filer, der allerede er behandlet. Silver-tabellerne bliver materialized views, som Databricks automatisk opdaterer. SQL-koden er kort og tydelig.

Materialized views understøtter ikke slowly changing dimensions direkte. SCD2 kræver derfor en separat MERGE-notebook, som bryder det deklarative mønster. Det er en enkel løsning, når SCD2 ikke er nødvendigt. Ellers må pipelinen blande to tilgange.

4. dbt-core på Databricks

Kode: src/dbt_project

dbt organiserer arbejdet omkring modeller: SQL SELECT-statements, som definerer, hvad en tabel skal indeholde. dbt håndterer rækkefølgen af dependencies, tests og dokumentation.

SCD2 håndteres via dbt snapshots, som sammenligner rækketilstande mellem kørsler:

{`{% snapshot snap_dim_customer %}
{{ config(strategy='check', unique_key='customer_id',
          check_cols=['email','address','city','country','segment']) }}
SELECT * FROM {{ ref('silver_customers') }}
{% endsnapshot %}`}

Workflowet er lidt usædvanligt. Fordi dbt snapshots fanger ændringer mellem kørsler (ikke mellem batches), kører bundlet en to-faset proces: først indlæs batch 1, snapshot, derefter indlæs batch 2, snapshot igen, og så byg gold. valid_from / valid_to-værdierne udledes fra _batch_id i stedet for snapshot-timestamps, da vi har brug for deterministiske forretningsdatoer, ikke kørselstidspunkter.

Den største fordel ved dbt er strukturen. Værktøjet håndhæver en projektstruktur, ref-baserede dependencies og en testkultur med not_null, unique og relationships. Hvis dit team allerede bruger dbt, eller I lægger vægt på portable SQL-kompetencer, er det et stærkt valg.

Strukturen betyder også endnu et værktøj at vedligeholde, blandt andet dbt_project.yml, profiles.yml og schema.yml. dbt's model passer heller ikke perfekt til streaming og pipeline-native features i Databricks. dbt-databricks-adapteren er et godt startpunkt, hvis du vil prøve det.

5. Delta Live Tables: den klassiske syntaks

Kode: pipeline.sql

Delta Live Tables (DLT) var Databricks' første rigtige forsøg på deklarative pipelines. Det introducerede nogle genuint nyttige idéer: indbyggede data quality expectations, nativ SCD2 via APPLY CHANGES og en managed runtime, der håndterer inkrementel behandling for dig.

CREATE STREAMING LIVE TABLE bronze_customers AS
SELECT * FROM STREAM read_files('/Volumes/.../customers', format => 'csv');

APPLY CHANGES INTO LIVE.dim_customer
FROM STREAM(LIVE.bronze_customers)
KEYS (customer_id)
STORED AS SCD TYPE 2;

SCD2 er højdepunktet. To linjer erstatter snesevis af linjer med manuel merge-logik. Runtimen producerer automatisk __START_AT / __END_AT-kolonner (lidt anderledes navngivning end valid_from / valid_to-konventionen, men semantisk ækvivalent).

Databricks har siden omdøbt produktet og moderniseret syntaksen. DLT-nøgleord som CREATE STREAMING LIVE TABLE og APPLY CHANGES INTO virker stadig, men er legacy-syntaks. Nye projekter bør bruge de aktuelle navne.

6. Lakeflow pipelines: den anbefalede tilgang

Kode: SQL-version, Python-version

Produktet, der tidligere hed DLT, kaldes nu Lakeflow pipelines. Det bruger samme runtime og har de samme muligheder, men syntaksen følger nu almindelige SQL-konventioner:

CREATE OR REFRESH STREAMING TABLE bronze_customers AS
SELECT * FROM STREAM read_files('/Volumes/.../customers', format => 'csv');

CREATE FLOW scd2_dim_customer AS AUTO CDC INTO dim_customer
FROM STREAM(bronze_customers)
KEYS (customer_id)
STORED AS SCD TYPE 2;

SQL-versionen beskriver pipelinen direkte. Python-versionen i repoet bruger det kompatible dlt-modul med @dlt.table og create_auto_cdc_flow(). Til ny kode anbefaler Databricks nu from pyspark import pipelines as dp og @dp.table.

Lakeflow har flere praktiske fordele:

  • Inkrementel som standard. Streaming tables holder styr på, hvad der er behandlet. Du skriver ikke checkpoint-logik.
  • Nativ SCD2. AUTO CDC INTO håndterer historikken uden en manuel merge eller snapshot-workaround.
  • Indbygget datakvalitet. CONSTRAINT ... EXPECT ... ON VIOLATION DROP ROW validerer data inline.
  • Én pipeline-definition. Hele Bronze-Silver-Gold-flowet ligger i én fil. Databricks håndterer afhængighedsopløsning og udførelsesrækkefølge.
  • Mindre kode samlet set. SQL-pipelinen er én fil, der erstatter tre notebooks med manuel logik.

Pipeline-syntaksen har sin egen semantik, som du skal lære. Debugging foregår i pipeline-runtimen og ikke i en notebook, du kan steppe igennem. Kolonnenavnene __START_AT og __END_AT adskiller sig også fra valid_from og valid_to, hvilket er en lille forskel, som er værd at dokumentere for teamet.

For de fleste teams, der bygger nye medallion-pipelines på Databricks, er det her, jeg ville starte.

Så hvilken skal du bruge?

Mit valg ville afhænge af teamet og projektet:

Start med Lakeflow pipelines i et nyt Databricks-projekt. De balancerer læsbarhed, indbygget SCD2-understøttelse, data quality checks og inkrementel behandling med mindst mulig custom kode.

Brug dbt, hvis dit team allerede er investeret i dbt-økosystemet, eller hvis portabilitet på tværs af platforme er vigtig. dbt's test- og dokumentationskultur er en klar fordel, selv om den giver mere projektarbejde.

Brug Python eller SQL notebooks til prototyper, eksperimenter, eller når du har brug for fuld kontrol over behandlingslogikken. De er hurtige at starte og nemme at forstå, men de akkumulerer hurtigt kompleksitet i produktion.

Brug Materialized Views og Streaming Tables til kompakte SQL-only pipelines, hvor SCD2 ikke er nødvendigt. I den situation er det den enkleste løsning.

Undgå klassisk DLT-syntaks til nye projekter. Den kører stadig fint, men der er ingen grund til at starte med forældede nøgleord, når de moderne ækvivalenter gør det samme.

Den fulde kildekode ligger på github.com/romaklimenko/databricks-medallion. Hvis du vil dykke ned i en specifik tilgang, har README'en detaljerede noter om hver enkelt.