Pipeline di produzione con i workflow

Trasformazione dei dati con Spark SQL in Databricks

Disha Mukherjee

Lead Data Engineer

Perché Delta Lake?

Un caveau digitale luminoso e sicuro con tabelle dati ordinate e uno scudo protettivo, stile moderno piatto

 

$$

  • Transazioni ACID → eseguono rollback sui write falliti
  • Enforcement dello schema → blocca i tipi non corrispondenti
  • Versioning → interroga qualsiasi stato precedente
Trasformazione dei dati con Spark SQL in Databricks

Scrivere su Delta

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

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

$$

$$

  • In Unity Catalog compare una nuova tabella Delta
Trasformazione dei dati con Spark SQL in Databricks

Attività del notebook

task1_ingest: carica e pulisci

  • Carica CSV → applica pulizia → scrive una tabella Delta
Trasformazione dei dati con Spark SQL in Databricks

Attività del notebook

task2_metrics: ricavi per categoria

  • Legge la tabella pulita → calcola metriche → scrive in una nuova tabella
Trasformazione dei dati con Spark SQL in Databricks

Attività del notebook

task3_customers: classifica per spesa

  • Legge la tabella pulita → classifica i clienti → salva in una nuova tabella
Trasformazione dei dati con Spark SQL in Databricks

Creare il job

Interfaccia Jobs and Pipelines che mostra un job con tre task completati in un DAG, con frecce di dipendenza da task1_ingest a task2_metrics a task3_customers

  • Puoi eseguire manualmente, pianificare job e impostare trigger
Trasformazione dei dati con Spark SQL in Databricks

Esecuzione del job

DAG dell'esecuzione del job con task1_ingest e task2_metrics riusciti in verde e task3_customers fallito in rosso

Trasformazione dei dati con Spark SQL in Databricks

Esecuzione del job

Controllo di un errore nel terzo task

  • Errore: customer-id dovrebbe essere Customer_ID
Trasformazione dei dati con Spark SQL in Databricks

Esecuzione del job

Viste grafiche con successo

Una vista timeline

Trasformazione dei dati con Spark SQL in Databricks

Cos'è Lakeflow?

$$

$$

confronto: Imperativo (Jobs) | Dichiarativo (Lakeflow)

 

  • Jobs → gestiamo ogni passaggio
  • Lakeflow → dichiari cosa devono contenere le tabelle
  • Databricks gestisce ordine, retry e compute
Trasformazione dei dati con Spark SQL in Databricks

Il pattern @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(...)
Trasformazione dei dati con Spark SQL in Databricks

Esecuzione della pipeline

DAG della pipeline Lakeflow che mostra transactions_bronze 100K righe poi transactions_silver 33K righe poi category_revenue_gold 6 righe, tutte viste materializzate verdi

  • Bronze (100K righe) → Silver (33K righe) → Gold (6 righe)
Trasformazione dei dati con Spark SQL in Databricks

Notebooks, Jobs o Lakeflow?

$$

layer: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → esplorazione, prototipi
  • Databricks Jobs → pipeline multi-step pianificate
  • Lakeflow → completamente gestito, dichiarativo
Trasformazione dei dati con Spark SQL in Databricks

Esercitiamoci!

Trasformazione dei dati con Spark SQL in Databricks

Preparing Video For Download...