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 的正确使用、集群的精细控制和作业化的工程思维。