Производственные пайплайны с 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

Создание задания

Интерфейс Jobs and Pipelines с завершённым DAG из трёх задач: стрелки зависимостей от task1_ingest к task2_metrics и task3_customers

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

Запуск задания

DAG запуска задания: task1_ingest и task2_metrics выполнены успешно (зелёный), task3_customers завершился с ошибкой (красный)

Преобразование данных с помощью Spark SQL в Databricks

Запуск задания

Проверка ошибки в третьей задаче

  • Ошибка: customer-id должен быть Customer_ID
Преобразование данных с помощью Spark SQL в Databricks

Запуск задания

Граф выполнения с успешным результатом

Вид временной шкалы

Преобразование данных с помощью 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 (100К строк) → transactions_silver (33К строк) → category_revenue_gold (6 строк), все представления материализованы и выделены зелёным

  • Bronze (100К строк) → Silver (33К строк) → 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...