高效聚合与连接数据

在 Databricks 中使用 Spark SQL 进行数据转换

Disha Mukherjee

Lead Data Engineer

三个问题,一个干净数据集

recraft: half: 一位数据分析师在多块屏幕上查看彩色图表和业务报告,透明背景独立呈现

 

$$

  • 哪些品类带来最高收入?
  • 谁是头部客户
  • 如何用部门信息丰富结果?
在 Databricks 中使用 Spark SQL 进行数据转换

groupBy() 与 agg() 的工作方式

 

$$

  • groupBy() —— 按键列分区行
  • agg() —— 并行对各组应用函数
  • 二者均为惰性执行——需动作触发计算
  • 一次 agg() 调用 = 一次数据扫描

recraft: half: 彩色数据行被拆分为三个带标签的分组,每组下方计算汇总值,透明背景独立呈现

在 Databricks 中使用 Spark SQL 进行数据转换

按品类收入——代码

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())
)
在 Databricks 中使用 Spark SQL 进行数据转换

按品类收入——输出

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       |
+-----------+-----------------+-----------------+---------------+
在 Databricks 中使用 Spark SQL 进行数据转换

用维度表进行丰富

recraft: half: 两张数据库表通过发光的链接或桥相连,透明背景独立呈现

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
在 Databricks 中使用 Spark SQL 进行数据转换

标准左连接

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          |
+-----------+-----------+----------+------------------+
在 Databricks 中使用 Spark SQL 进行数据转换

隐藏成本——shuffle

recraft: half: 数据包在由服务器与节点组成的网络间流动,箭头指示方向,透明背景独立呈现

 

  • 匹配键须落在同一台机器
  • Spark 跨网络传输行以对齐
  • Shuffle 在大规模下是瓶颈
在 Databricks 中使用 Spark SQL 进行数据转换

Spark UI

recraft: full: 深色分析监控面板,含彩色横向条形图显示作业阶段、进度条和数据吞吐指标,透明背景独立呈现

 

  • JobsStages —— 跟踪执行与耗时
  • Shuffle Read/Write —— 观察数据搬移量
  • 有助于定位流水线瓶颈
在 Databricks 中使用 Spark SQL 进行数据转换

.explain() 读取查询计划

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() —— 检查执行计划
  • 节点 15 与 16 —— 记录了路由参数
在 Databricks 中使用 Spark SQL 进行数据转换

广播连接——优化方案

# 用 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
在 Databricks 中使用 Spark SQL 进行数据转换

前后对比——验证计划

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 无需搬移
  • 一处改动,规模化可从分钟降到秒
在 Databricks 中使用 Spark SQL 进行数据转换

让我们一起练习吧!

在 Databricks 中使用 Spark SQL 进行数据转换

Preparing Video For Download...