Produktionspipelines med workflows

Datatransformation med Spark SQL i Databricks

Disha Mukherjee

Lead Data Engineer

Varför Delta Lake?

Ett säkert, glödande digitalt valv med organiserade datatabeller och en skyddssköld, modernt platt stilformat

 

$$

  • ACID-transaktioner → återställer misslyckade skrivningar
  • Schematvång → blockerar typmissmatchningar
  • Versionshantering → fråga vilket tidigare tillstånd som helst
Datatransformation med Spark SQL i Databricks

Skriva till Delta

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • En ny Delta-tabell visas i Unity Catalog
Datatransformation med Spark SQL i Databricks

Notebook-uppgifter

task1_ingest: läs in och rensa

  • Läser in CSV → rensar data → skriver en Delta-tabell
Datatransformation med Spark SQL i Databricks

Notebook-uppgifter

task2_metrics: intäkter per kategori

  • Läser rensat register → beräknar mätvärden → skriver till ny tabell
Datatransformation med Spark SQL i Databricks

Notebook-uppgifter

task3_customers: rangordna efter utgifter

  • Läser rensat register → rangordnar kunder → sparar till ny tabell
Datatransformation med Spark SQL i Databricks

Skapa jobbet

Gränssnittet Jobs and Pipelines visar ett avslutat jobb-DAG med tre uppgifter och beroendepilarna från task1_ingest till task2_metrics till task3_customers

  • Kan köras manuellt, schemaläggas och utlösas via triggers
Datatransformation med Spark SQL i Databricks

Köra jobbet

Jobb-DAG med task1_ingest och task2_metrics markerade gröna som slutförda och task3_customers markerad röd som misslyckad

Datatransformation med Spark SQL i Databricks

Köra jobbet

Kontrollera ett fel i uppgift tre

  • Fel: customer-id ska vara Customer_ID
Datatransformation med Spark SQL i Databricks

Köra jobbet

Grafvy med lyckad körning

En tidslinjevy

Datatransformation med Spark SQL i Databricks

Vad är Lakeflow?

$$

$$

jämförelse: Imperativ (Jobs) | Deklarativ (Lakeflow)

 

  • Jobs → vi hanterar varje steg
  • Lakeflow → deklarera vad tabeller ska innehålla
  • Databricks hanterar ordning, omförsök och beräkning
Datatransformation med Spark SQL i Databricks

Mönstret @dlt.table

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
Datatransformation med Spark SQL i Databricks

Pipelinekörning

Lakeflow pipeline-DAG som visar transactions_bronze med 100 000 rader, sedan transactions_silver med 33 000 rader, sedan category_revenue_gold med 6 rader – alla gröna materialiserade vyer

  • Bronze (100 000 rader) → Silver (33 000 rader) → Gold (6 rader)
Datatransformation med Spark SQL i Databricks

Notebooks, Jobs eller Lakeflow?

$$

lager: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → utforskning och prototyputveckling
  • Databricks Jobs → flerstegs-, schemalagda pipelines
  • Lakeflow → helt hanterad, deklarativ
Datatransformation med Spark SQL i Databricks

Nu kör vi en övning!

Datatransformation med Spark SQL i Databricks

Preparing Video For Download...