Databricks 性能优化实战

早八人AI炼丹师 专家 2026/7/27 540 浏览 1 点赞 约 1 分钟

盲目增加集群 Worker 节点根本解决不了 Spark 的性能问题,因为大部分卡顿都源于 Shuffle 和数据倾斜(Skew)。如果 Join 的分区分布极不均匀,堆再多机器也只是让大部分节点在干等那个最慢的任务。

分享一个基于 Delta Lake 和 Unity Catalog 的实战优化流程,涵盖从数据治理到查询加速的完整链路。

一、 Unity Catalog 基础环境搭建

为了实现统一的权限管理和血缘追踪,所有表必须注册在 Unity Catalog 中。
CREATE CATALOG IF NOT EXISTS retail_analytics;
CREATE SCHEMA IF NOT EXISTS retail_analytics.events;

CREATE TABLE IF NOT EXISTS retail_analytics.events.raw_orders (
 order_id STRING,
 customer_id STRING,
 product_id STRING,
 quantity INT,
 event_ts TIMESTAMP
)
USING DELTA;

CREATE TABLE IF NOT EXISTS retail_analytics.events.dim_products (
 product_id STRING,
 category STRING,
 unit_cost DOUBLE
)
USING DELTA;

GRANT SELECT ON TABLE retail_analytics.events.raw_orders TO `analysts`;
Databricks 性能优化实战

二、 针对小表强制执行 Broadcast Join

当维度表较小时,通过 broadcast 提示可以避免全量 Shuffle。虽然 Spark 的 CBO 优化器通常会自动处理,但手动指定可以防止维度表在增长到 spark.sql.autoBroadcastJoinThreshold(默认 10MB)阈值以上时突然导致性能崩塌。
from pyspark.sql import functions as F

orders = spark.table("retail_analytics.events.raw_orders")
products = spark.table("retail_analytics.events.dim_products")

# 强制广播小表,避免大表 shuffle
joined = orders.join(F.broadcast(products), on="product_id", how="left")

三、 识别与处理数据倾斜

在执行 groupBy("customer_id") 这种宽依赖转换时,如果某个 customer_id(如
from pyspark.sql import functions as F

SALT_BUCKETS = 20

# 给 Key 增加随机后缀,将大 Key 分散到多个分区
df_salted = joined.withColumn(
    "salted_customer_id", 
    F.concat(F.col("customer_id"), F.lit("_"), F.expr(f"floor(rand() * {SALT_BUCKETS})"))
)
B 大客户)的数据量远超其他用户,就会出现一个 Task 运行几十分钟,而其他 Task 秒完的情况。

虽然现代 Databricks Runtime 的 AQE(自适应查询执行)能处理部分倾斜,但对于极端的 Hot Key,最稳妥的实操方案依然是加盐(Salting)

from pyspark.sql import functions as F

SALT_BUCKETS = 20

# 给 Key 增加随机后缀,将大 Key 分散到多个分区
df_salted = joined.withColumn(
    "salted_customer_id", 
    F.concat(F.col("customer_id"), F.lit("_"), F.expr(f"floor(rand() * {SALT_BUCKETS})"))
)

四、 Delta Lake 的 Z-Ordering 优化

最后一步是物理存储优化。对于下游经常需要过滤的列,使用 Z-Ordering 重新排列数据,可以让 Spark 在读取时跳过不相关的 Parquet 文件,极大提升查询速度。
-- 对经常出现在 WHERE 子句中的列进行 Z-Ordering
OPTIMIZE retail_analytics.events.aggregated_orders 
ZORDER BY (customer_id, category);

核心优化点总结:

  • 维度表:broadcast 消除 Shuffle。
  • 计算倾斜: 开启 AQE 或手动对 Key 加盐。
  • 存储层: 使用 OPTIMIZE + ZORDER 实现文件跳过。
Databricks 性能优化实战

AI编程AI编程实战pythondatabricksspark

全部回复 (4)

前端大鹏 初级 2026/7/27

没开AQE就去手动调分区数,简直是自虐行为!

0 回复
数据分析师大山 中级 2026/7/27

加盐大法简直是救星,比死磕内存配置快太多了。

0 回复
程序员Tom 高级 2026/7/27

傻傻加机器简直是浪费钱,早点调倾斜能省多少资源!

0 回复
远程办公产品狗 中级 2026/7/27

数据倾斜这坑太深了,快分享下是用哪个监控指标揪出问题的!

0 回复

发表回复

支持 Markdown 格式