Potoki produkcyjne z workflows

Data Transformation with Spark SQL in Databricks

Disha Mukherjee

Lead Data Engineer

Dlaczego Delta Lake?

Bezpieczny, świecący cyfrowy sejf z uporządkowanymi tabelami danych i tarczą ochronną, płaski nowoczesny styl

 

$$

  • Transakcje ACID → cofanie nieudanych zapisów
  • Wymuszanie schematu → blokowanie niezgodnych typów
  • Wersjonowanie → odpytywanie dowolnego wcześniejszego stanu
Data Transformation with Spark SQL in Databricks

Zapis do Delta

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

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

$$

$$

  • W Unity Catalog pojawia się nowa tabela Delta
Data Transformation with Spark SQL in Databricks

Zadania notebooka

task1_ingest: wczytaj i oczyść dane

  • Wczytuje CSV → czyści dane → zapisuje tabelę Delta
Data Transformation with Spark SQL in Databricks

Zadania notebooka

task2_metrics: przychody według kategorii

  • Odczytuje oczyszczoną tabelę → oblicza metryki → zapisuje do nowej tabeli
Data Transformation with Spark SQL in Databricks

Zadania notebooka

task3_customers: rankowanie według wydatków

  • Odczytuje oczyszczoną tabelę → rankuje klientów → zapisuje do nowej tabeli
Data Transformation with Spark SQL in Databricks

Tworzenie zadania

Interfejs Jobs and Pipelines z ukończonym grafem DAG trzech zadań: strzałki zależności od task1_ingest do task2_metrics do task3_customers

  • Możliwe uruchomienie ręczne, harmonogramowanie i wyzwalacze
Data Transformation with Spark SQL in Databricks

Uruchamianie zadania

Graf przebiegu zadania: task1_ingest i task2_metrics zakończone sukcesem (zielone), task3_customers zakończone błędem (czerwone)

Data Transformation with Spark SQL in Databricks

Uruchamianie zadania

Analiza błędu w zadaniu trzecim

  • Błąd: customer-id powinno być Customer_ID
Data Transformation with Spark SQL in Databricks

Uruchamianie zadania

Widok grafu — wszystkie zadania zakończone sukcesem

Widok osi czasu

Data Transformation with Spark SQL in Databricks

Czym jest Lakeflow?

$$

$$

porównanie: Imperatywne (Jobs) | Deklaratywne (Lakeflow)

 

  • Jobs → sami zarządzamy każdym krokiem
  • Lakeflow → deklarujemy, co tabele mają zawierać
  • Databricks obsługuje kolejność, ponowne próby i obliczenia
Data Transformation with Spark SQL in Databricks

Wzorzec @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(...)
Data Transformation with Spark SQL in Databricks

Przebieg potoku

Graf DAG potoku Lakeflow: transactions_bronze (100 tys. wierszy) → transactions_silver (33 tys. wierszy) → category_revenue_gold (6 wierszy), wszystkie zmaterializowane widoki w kolorze zielonym

  • Bronze (100 tys. wierszy) → Silver (33 tys. wierszy) → Gold (6 wierszy)
Data Transformation with Spark SQL in Databricks

Notebooks, Jobs czy Lakeflow?

$$

warstwy: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → eksploracja i prototypowanie
  • Databricks Jobs → wieloetapowe, harmonogramowane potoki
  • Lakeflow → w pełni zarządzane, deklaratywne
Data Transformation with Spark SQL in Databricks

Czas na praktykę!

Data Transformation with Spark SQL in Databricks

Preparing Video For Download...