Pipeline-uri de producție cu fluxuri de lucru

Transformarea datelor cu Spark SQL în Databricks

Disha Mukherjee

Lead Data Engineer

De ce Delta Lake?

Un seif digital securizat, luminos, cu tabele de date organizate și un scut de protecție, stil modern plat

 

$$

  • Tranzacții ACID → anulează scrierile eșuate
  • Aplicarea schemei → blochează tipurile incompatibile
  • Versionare → interogează orice stare anterioară
Transformarea datelor cu Spark SQL în Databricks

Scriere în Delta

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

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

$$

$$

  • Un tabel Delta nou apare în Unity Catalog
Transformarea datelor cu Spark SQL în Databricks

Sarcini notebook

task1_ingest: încărcare și curățare

  • Încarcă CSV → aplică curățarea → scrie un tabel Delta
Transformarea datelor cu Spark SQL în Databricks

Sarcini notebook

task2_metrics: venituri pe categorie

  • Citește tabelul curățat → calculează metrici → scrie într-un tabel nou
Transformarea datelor cu Spark SQL în Databricks

Sarcini notebook

task3_customers: clasificare după cheltuieli

  • Citește tabelul curățat → clasifică clienții → salvează într-un tabel nou
Transformarea datelor cu Spark SQL în Databricks

Crearea job-ului

Interfața Jobs and Pipelines cu un DAG finalizat din trei sarcini, cu săgeți de dependență de la task1_ingest la task2_metrics la task3_customers

  • Poate rula manual, programa job-uri și seta declanșatoare
Transformarea datelor cu Spark SQL în Databricks

Rularea job-ului

DAG al rulării job-ului cu task1_ingest și task2_metrics reușite în verde și task3_customers eșuat în roșu

Transformarea datelor cu Spark SQL în Databricks

Rularea job-ului

Verificarea unei erori în sarcina trei

  • Eroare: customer-id ar trebui să fie Customer_ID
Transformarea datelor cu Spark SQL în Databricks

Rularea job-ului

Vizualizare grafică cu succes

O vizualizare cronologică

Transformarea datelor cu Spark SQL în Databricks

Ce este Lakeflow?

$$

$$

comparație: Imperativ (Jobs) | Declarativ (Lakeflow)

 

  • Jobs → gestionăm fiecare pas
  • Lakeflow → declarăm ce trebuie să conțină tabelele
  • Databricks gestionează ordinea, reluările și calculul
Transformarea datelor cu Spark SQL în Databricks

Șablonul @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(...)
Transformarea datelor cu Spark SQL în Databricks

Rularea pipeline-ului

DAG-ul pipeline-ului Lakeflow cu transactions_bronze 100K rânduri, transactions_silver 33K rânduri și category_revenue_gold 6 rânduri, toate vizualizări materializate verzi

  • Bronze (100K rânduri) → Silver (33K rânduri) → Gold (6 rânduri)
Transformarea datelor cu Spark SQL în Databricks

Notebooks, Jobs sau Lakeflow?

$$

straturi: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → explorare, prototipare
  • Databricks Jobs → pipeline-uri multi-pas, programate
  • Lakeflow → complet gestionat, declarativ
Transformarea datelor cu Spark SQL în Databricks

Hai să exersăm!

Transformarea datelor cu Spark SQL în Databricks

Preparing Video For Download...