Produkční pipelines s workflows

Transformace dat pomocí Spark SQL v Databricks

Disha Mukherjee

Lead Data Engineer

Proč Delta Lake?

Zabezpečený zářící digitální trezor s uspořádanými datovými tabulkami a ochranným štítem, plochý moderní styl

 

$$

  • ACID transakce → vrácení neúspěšných zápisů zpět
  • Vynucení schématu → blokování neodpovídajících typů
  • Verzování → dotaz na libovolný předchozí stav
Transformace dat pomocí Spark SQL v Databricks

Zápis do Delty

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

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

$$

$$

  • V Unity Catalogu se objeví nová Delta tabulka
Transformace dat pomocí Spark SQL v Databricks

Notebookové tasky

task1_ingest: načtení a čištění

  • Načte CSV → provede čištění → zapíše Delta tabulku
Transformace dat pomocí Spark SQL v Databricks

Notebookové tasky

task2_metrics: tržby podle kategorie

  • Přečte vyčištěnou tabulku → vypočítá metriky → zapíše do nové tabulky
Transformace dat pomocí Spark SQL v Databricks

Notebookové tasky

task3_customers: řazení podle útraty

  • Přečte vyčištěnou tabulku → seřadí zákazníky → uloží do nové tabulky
Transformace dat pomocí Spark SQL v Databricks

Vytvoření jobu

Rozhraní Jobs and Pipelines zobrazující dokončený DAG jobu se třemi tasky a šipkami závislostí od task1_ingest přes task2_metrics po task3_customers

  • Lze spustit ručně, naplánovat nebo nastavit triggery
Transformace dat pomocí Spark SQL v Databricks

Spuštění jobu

DAG spuštění jobu: task1_ingest a task2_metrics úspěšné zeleně, task3_customers selhalo červeně

Transformace dat pomocí Spark SQL v Databricks

Spuštění jobu

Kontrola chyby ve třetím tasku

  • Chyba: customer-id má být Customer_ID
Transformace dat pomocí Spark SQL v Databricks

Spuštění jobu

Grafové zobrazení s úspěchem

Časové zobrazení

Transformace dat pomocí Spark SQL v Databricks

Co je Lakeflow?

$$

$$

srovnání: Imperativní (Jobs) | Deklarativní (Lakeflow)

 

  • Jobs → každý krok spravujeme sami
  • Lakeflow → deklarujeme, co mají tabulky obsahovat
  • Databricks řídí pořadí, opakování i výpočty
Transformace dat pomocí Spark SQL v Databricks

Vzor @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(...)
Transformace dat pomocí Spark SQL v Databricks

Spuštění pipeline

DAG Lakeflow pipeline zobrazující transactions_bronze 100 000 řádků, pak transactions_silver 33 000 řádků, pak category_revenue_gold 6 řádků – vše zelené materializované pohledy

  • Bronze (100 000 řádků) → Silver (33 000 řádků) → Gold (6 řádků)
Transformace dat pomocí Spark SQL v Databricks

Notebooks, Jobs, nebo Lakeflow?

$$

vrstvy: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → průzkum a prototypování
  • Databricks Jobs → vícekrokové, naplánované pipelines
  • Lakeflow → plně spravované, deklarativní
Transformace dat pomocí Spark SQL v Databricks

Pojďme cvičit!

Transformace dat pomocí Spark SQL v Databricks

Preparing Video For Download...