Промислові конвеєри з workflows

Перетворення даних за допомогою Spark SQL у Databricks

Disha Mukherjee

Lead Data Engineer

Чому Delta Lake?

Надійне сяйливе цифрове сховище з упорядкованими таблицями даних і захисним щитом, плаский сучасний стиль

 

$$

  • ACID-транзакції → відкат невдалих записів
  • Контроль схеми → блокування невідповідних типів
  • Версіонування → запити до будь-якого попереднього стану
Перетворення даних за допомогою Spark SQL у Databricks

Запис у Delta

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

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

$$

$$

  • У Unity Catalog з'являється нова таблиця Delta
Перетворення даних за допомогою Spark SQL у Databricks

Завдання в ноутбуках

task1_ingest: завантаження й очищення

  • Завантажує CSV → очищує → записує таблицю Delta
Перетворення даних за допомогою Spark SQL у Databricks

Завдання в ноутбуках

task2_metrics: виручка за категоріями

  • Читає очищену таблицю → обчислює метрики → записує в нову таблицю
Перетворення даних за допомогою Spark SQL у Databricks

Завдання в ноутбуках

task3_customers: рейтинг за витратами

  • Читає очищену таблицю → ранжує клієнтів → зберігає в нову таблицю
Перетворення даних за допомогою Spark SQL у Databricks

Створення job

Інтерфейс Jobs and Pipelines із завершеним DAG на три завдання та стрілками залежностей від task1_ingest до task2_metrics до task3_customers

  • Можна запускати вручну, планувати й налаштовувати тригери
Перетворення даних за допомогою Spark SQL у Databricks

Запуск job

DAG запуску job: task1_ingest і task2_metrics успішні (зелений), task3_customers помилка (червоний)

Перетворення даних за допомогою Spark SQL у Databricks

Запуск job

Перевірка помилки в третьому завданні

  • Помилка: customer-id має бути Customer_ID
Перетворення даних за допомогою Spark SQL у Databricks

Запуск job

Графіки з успішним результатом

Часова шкала

Перетворення даних за допомогою Spark SQL у Databricks

Що таке Lakeflow?

$$

$$

порівняння: Імперативний (Jobs) | Декларативний (Lakeflow)

 

  • Jobs → ми керуємо кожним кроком
  • Lakeflow → оголошуємо, що мають містити таблиці
  • Databricks визначає порядок, повтори й обчислення
Перетворення даних за допомогою Spark SQL у Databricks

Шаблон @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(...)
Перетворення даних за допомогою Spark SQL у Databricks

Запуск конвеєра

DAG конвеєра Lakeflow: transactions_bronze 100K рядків, далі transactions_silver 33K рядків, далі category_revenue_gold 6 рядків — усе зелені матеріалізовані подання

  • Bronze (100K рядків) → Silver (33K) → Gold (6)
Перетворення даних за допомогою Spark SQL у Databricks

Notebooks, Jobs чи Lakeflow?

$$

шари: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → дослідження, прототипування
  • Databricks Jobs → багатокрокові, заплановані конвеєри
  • Lakeflow → повністю керований, декларативний
Перетворення даних за допомогою Spark SQL у Databricks

Давайте потренуємось!

Перетворення даних за допомогою Spark SQL у Databricks

Preparing Video For Download...