Агрегация и объединение данных

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

Disha Mukherjee

Lead Data Engineer

Три вопроса — один чистый набор данных

recraft: half: Аналитик данных изучает красочные графики и бизнес-отчёты на нескольких экранах, изолированный на прозрачном фоне

 

$$

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

Как работают groupBy() и agg()

 

$$

  • groupBy() — разбивает строки по ключевому столбцу
  • agg() — применяет функции к каждой группе параллельно
  • Оба ленивые — вычисление запускается действием
  • Один вызов agg() = одно сканирование данных

recraft: half: Разноцветные строки данных, разбитые на три отдельные группы с меткой и итоговым значением под каждой группой, изолированные на прозрачном фоне

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

Выручка по категориям — код

category_revenue = (
    df_valid
    .groupBy("Category")

.agg( F.round(F.sum("Transaction_Amount"), 2).alias("total_revenue"), F.count("Transaction_Amount").alias("transaction_count"), F.round(F.avg("Transaction_Amount"), 2).alias("avg_transaction"), )
.orderBy(F.col("total_revenue").desc())
)
Преобразование данных с помощью Spark SQL в Databricks

Выручка по категориям — вывод

category_revenue.show(truncate=False)
+-----------+-----------------+-----------------+---------------+
|Category   |total_revenue    |transaction_count|avg_transaction|
+-----------+-----------------+-----------------+---------------+
|Clothing   |338266605.24     |6742             |50173.04       |
|Dining     |337631718.64     |6719             |50250.29       |
|Electronics|331102843.72     |6614             |50060.91       |
|Savings    |329346609.07     |6596             |49931.26       |
|Groceries  |329003451.57     |6537             |50329.43       |
|Unknown    |880591.27        |15               |58706.08       |
+-----------+-----------------+-----------------+---------------+
Преобразование данных с помощью Spark SQL в Databricks

Обогащение данных таблицей измерений

recraft: half: Две таблицы базы данных, соединённые светящейся связью или мостом, изолированные на прозрачном фоне

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Преобразование данных с помощью Spark SQL в Databricks

Стандартное левое соединение

df_joined = df_valid.join(df_dim, on="Category", how="left")

df_joined.select( "Customer_ID", "Category", "Department", "Transaction_Amount" ).show(5, truncate=False)
+-----------+-----------+----------+------------------+
|Customer_ID|Category   |Department|Transaction_Amount|
+-----------+-----------+----------+------------------+
|CUST003    |Electronics|Tech      |5752.36           |
|CUST009    |Clothing   |Retail    |28959.12          |
|CUST010    |Savings    |Finance   |72098.18          |
|CUST011    |Savings    |Finance   |49771.05          |
|CUST020    |Groceries  |Food      |69825.14          |
+-----------+-----------+----------+------------------+
Преобразование данных с помощью Spark SQL в Databricks

Скрытые издержки — перемешивание данных

recraft: half: Потоки пакетов данных по сети связанных серверов и узлов со стрелками, показывающими движение, изолированные на прозрачном фоне

 

  • Совпадающие ключи должны попасть на одну машину
  • Spark передаёт строки по сети для их выравнивания
  • Перемешивание = узкое место при масштабировании
Преобразование данных с помощью Spark SQL в Databricks

Spark UI

recraft: full: Тёмный аналитический дашборд с яркими горизонтальными столбчатыми диаграммами, отображающий этапы задач, индикаторы прогресса и метрики пропускной способности данных, изолированный на прозрачном фоне

 

  • Задачи и Этапы — отслеживают, что выполнялось и сколько времени
  • Shuffle Read/Write — показывает объём перемещённых данных
  • Помогает диагностировать узкие места в пайплайне
Преобразование данных с помощью Spark SQL в Databricks

Чтение плана запроса с помощью .explain()

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain()анализ плана выполнения
  • Узлы 15 и 16 — аргументы маршрутизации не записаны
Преобразование данных с помощью Spark SQL в Databricks

Широковещательное соединение — решение

# Wrap the small table in F.broadcast()
df_broadcast = df_valid.join(
    F.broadcast(df_dim),
    on="Category",
    how="left"
)

print(f"Broadcast joined rows: {df_broadcast.count():,}")
Broadcast joined rows: 33,223
Преобразование данных с помощью Spark SQL в Databricks

До и после — проверка плана

df_broadcast.explain(mode="formatted")
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...            
(15) PhotonShuffleExchangeSink
Arguments: SinglePartition
(16) PhotonShuffleMapStage
Arguments: EXECUTOR_BROADCAST, [id=#11839]
  • Таблица измерений получила инструкции маршрутизации
  • Большая таблица df_valid не перемещается
  • Одно изменение = минуты превращаются в секунды при масштабировании
Преобразование данных с помощью Spark SQL в Databricks

Давайте потренируемся!

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

Preparing Video For Download...