Databricks 进阶用法:从入门到精通

data-science进阶7 分钟阅读2026/9/26

Databricks 进阶用法:从入门到精通——我踩过的坑和总结的实战经验

一个真实的痛点

去年我们团队接了个活:每天要处理 300GB 左右的用户行为日志,跑一批复杂的聚合和 join。一开始我把脚本直接扔进 Databricks 的 Notebook 里跑,结果每次都要 4 个小时,偶尔还因为内存溢出直接挂掉。老板问“集群那么贵,为什么跑这么慢”,我才发现自己只是把 Databricks 当成了一个“能跑 PySpark 的网页版 Jupyter”——这完全浪费了它的能力。

这篇文章就是我这一年从“会用”到“用好”的过程,希望你能少走弯路。

第一步:搞懂集群配置,别再开“坦克买菜”

我犯的第一个错误就是无脑开大集群。有一次我开了 8 个 Standard_DS4_v2 节点跑一个 5GB 的表聚合,结果 3 分钟跑完,账单让我肉疼。

后来我学到的核心原则:

  • 小数据(<10GB)用单机模式。Databricks 支持开启单节点集群(Single Node),没有 Spark worker,延迟低、成本低,跑小规模 ETL 和交互式探索非常合适。
  • 自动伸缩别设太宽。我的经验是 min/max 节点数控制在 2 倍以内,比如 min=2, max=4。设成 2 到 16 的话,autoscaling 反应慢,经常出现“任务已经挤死了,节点还在启动中”的尴尬。
  • 用 Spot 节点跑容错任务。非交互式的批处理可以勾选 Spot Instances,成本能降 60% 以上。注意:driver 节点别用 Spot,一个节点被回收整条作业就挂了。

另外一个容易忽略的点:打开自动终止(Auto Termination)。我有个同事半夜跑完作业忘了关集群,一个 Standard 集群空转了一整晚。

第二步:Delta Lake 才是精髓,别再用 Parquet 直读直写

刚上手时我把数据直接写成 Parquet,直到遇到两个问题:有人覆盖写坏了一张表;另一个下游任务读到了写到一半的脏数据。

换到 Delta Lake 之后,这两个问题都消失了:

# 直接把 Parquet 目录转成 Delta
df = spark.read.parquet("/mnt/raw/user_events")
df.write.format("delta").mode("overwrite").save("/mnt/delta/user_events")

几个我认为必须掌握的进阶操作:

1. MERGE 实现增量更新(替代全量重写)

from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, "/mnt/delta/user_events")
delta_table.alias("t").merge(
    updates_df.alias("s"),
    "t.user_id = s.user_id AND t.event_date = s.event_date"
).whenMatchedUpdateAll() \
 .whenNotMatchedInsertAll() \
 .execute()

我们每天 300GB 的重写任务改成 MERGE 增量后,跑批时间从 4 小时降到 40 分钟。

2. Time Travel 救过我一次命

有一次新同事的错误代码覆盖了几天的数据。如果没有 Delta,我们只能从备份恢复。实际操作:

-- 查看历史
DESCRIBE HISTORY delta.`/mnt/delta/user_events`;

-- 直接读回昨天的版本
SELECT * FROM delta.`/mnt/delta/user_events` VERSION AS OF 42;

如果确认要回滚:RESTORE TABLE ... TO VERSION AS OF 42; 五分钟解决事故。

3. 定期 OPTIMIZE 和 VACUUM

MERGE 用多了会产生大量小文件,查询会越来越慢。我们配了一个每天跑的维护作业:

OPTIMIZE delta.`/mnt/delta/user_events` WHERE event_date >= '2024-01-01';
VACUUM delta.`/mnt/delta/user_events` RETAIN 168 HOURS;

注意 VACUUM 前留足保留期(上面是 7 天),否则 Time Travel 的历史快照会被物理删除。

第三步:性能调优的几个实招

分区不是越多越好。 我一开始按 user_id 分区,结果分出了几千万个小文件,调度开销比计算还大。正确做法是按低基数字段分区,比如日期 event_date,配合 ZORDER BY (user_id) 加速查询:

OPTIMIZE events ZORDER BY (user_id);

Broadcast Join 处理倾斜。 一次 join 一个 20 亿行的行为表和 300 万行的用户维表,跑了一个小时没出结果。加上:

from pyspark.sql.functions import broadcast
df = events.join(broadcast(users), "user_id")

提前广播小表,避免 shuffle,20 分钟跑完。如果 join key 有数据倾斜(某个大 V 的行为占了 10% 数据),可以考虑给倾斜 key 加随机前缀打散。

开启 Photon。 如果预算允许,Photon 引擎对 SQL 和 DataFrame API 类负载提速通常在 2~5 倍,很多情况下比单纯加节点划算。

缓存要有节制。 df.cache() 不是万能的,缓存放不下反而会频繁淘汰。我的原则:只有被多次 Action 消费的大中间结果才缓存,用完 unpersist()。

第四步:用 Job Workflow 替代手动跑 Notebook

早期的我天天手动点 Run Notebook,后来用 Jobs 把整条链路串起来:

  • Notebook A:抽取增量数据 → Notebook B:MERGE 更新 → Notebook C:数据质量校验 → 失败自动重试 2 次 + 发 Slack 告警
  • 配置 cluster: existing_job_cluster,避免每次启动等 5 分钟
  • 用 Task 依赖而不是把所有逻辑塞进一个巨型 Notebook,哪个环节挂了一目了然

数据质量校验我用的是简单粗暴的断言:

assert df.filter("user_id IS NULL").count() == 0, "发现空 user_id,终止下游写入"

最后:诚实的局限性说明

  • 成本不透明是最大痛点。DBU 计费 + 云资源费用,很容易超预算,建议配合 system.billing 表和集群标签做成本归因。
  • Notebook 不适合生产代码。真正核心的逻辑应该抽成 Python 包或 Databricks Asset Bundles 管理,Notebook 只做编排。
  • 自动伸缩不是银弹,突发大任务的资源规划还是要自己拍。
  • 小文件问题和数据倾斜问题,工具只能缓解,根本上还是要靠设计。

一句话总结:Databricks 入门一天就够,但真正发挥它的价值,靠的是 Delta Lake 的正确使用、集群的精细控制和作业化的工程思维。

相关 Agent

H

拥抱未来

一个用于共享、训练和部署机器学习模型和数据集的平台。

了解更多 →