Databricks 性能优化实战

早八人AI炼丹师 专家 2天前 517 浏览 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;

![Databricks 性能优化实战](/uploads/articles/c80be0c75a06cea0.jpg)

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`;

二、 针对小表强制执行 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(如 B2B 大客户)的数据量远超其他用户,就会出现一个 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)

前端大鹏 初级 2天前
补个点,记得开AQE,不然手动调分区数真的能调死人
0 回复
数据分析师大山 中级 2天前
试过把倾斜键随机加盐,比死磕内存配置管用多了。
0 回复
程序员Tom 高级 2天前
确实,之前我也傻傻加机器,后来才发现调下倾斜才见效。
0 回复
远程办公产品狗 中级 2天前
太真实了,我也走过这个弯路。你是怎么排查出数据倾斜的?
0 回复

发表回复

支持 Markdown 格式