上周三凌晨两点,我的手机响了。生产环境的 Databricks 作业挂了。
我打开日志一看,一个原本只需 15 分钟的 ETL 任务已经跑了 3 个小时,最后因为内存溢出(OOM)崩了。更要命的是,下游的 Delta 表因为写入中断,处于一个半死不活的脏数据状态,导致所有依赖这张表的报表全部报错。
那天晚上我一边喝着冷咖啡一边修数据,一边在心里列清单:这些问题我踩过不止一次,是时候把 Databricks 里那些最容易让人栽跟头的坑和对应的解决方案系统地整理出来了。
如果你正在用 Databricks,不管是 Azure 还是 AWS 上的,以下这些场景你大概率会碰到。
问题一:Delta 表的并发写入冲突与脏数据
这是我凌晨那次事故的核心问题。当时我们有三个流式作业同时往同一张 Delta 表写数据,其中一个作业挂了之后留下了未提交的事务文件(也就是常说的脏文件),导致后续读取时报错 CONCURRENT_APPEND,或者直接读不到最新数据。
我的解决步骤:
首先,清理掉未提交的事务文件。在 Notebook 里运行:
spark.sql("VACUUM my_catalog.my_schema.sales_transactions RETAIN 0 HOURS")
注意,Databricks 默认的安全保留时间是 168 小时(7 天),直接运行上面这句会报错。你需要先关闭安全检查:
SET spark.databricks.delta.retentionDurationCheck.enabled = false;
清理完脏数据后,我重新配置了并发写入策略。对于同时追加写入的场景,在写入前设置乐观并发控制的容忍级别:
spark.conf.set("spark.databricks.delta.commitInfo.userMetadata", "etl_job_v3")
spark.sql("""
ALTER TABLE my_catalog.my_schema.sales_transactions
SET TBLPROPERTIES (delta.concurrentAppendWrites = true)
""")
踩坑提醒: VACUUM RETAIN 0 在生产环境是危险操作,只在紧急恢复时使用。日常运维中,把保留时间设为 24 小时就足够了:
VACUUM my_catalog.my_schema.sales_transactions RETAIN 24 HOURS
问题二:小文件灾难(表扫描极慢)
去年我们有一个业务表,每天增量写入 50 万条数据。跑了三个月后,一个简单的 SELECT COUNT(*) 居然要耗时 4 分钟。我查了一下文件分布:
display(spark.sql("DESC DETAIL my_catalog.my_schema.user_events").select("numFiles", "sizeInBytes"))
结果触目惊心:47 万个文件,总大小才 12GB。平均每个文件 26KB。Spark 在读取时需要为每个文件生成一个 Task,调度的开销远大于实际计算的开销。
我的解决方案是两步走:
第一步,立即用 OPTIMIZE 压缩已有小文件:
OPTIMIZE my_catalog.my_schema.user_events;
如果表按日期分区,可以只优化热分区,节省时间:
OPTIMIZE my_catalog.my_schema.user_events WHERE event_date >= '2024-01-01';
运行后,47 万个文件合并成了 3200 个,查询时间从 4 分钟降到了 8 秒。
第二步,从源头解决。在 Auto Loader 的写入逻辑中,加入 maxFilesPerTrigger 和 maxBytesPerTrigger 控制触发频率,并在写入时指定文件大小:
(df.writeStream
.format("delta")
.option("maxFilesPerTrigger", "100")
.option("maxBytesPerTrigger", "64mb")
.outputMode("append")
.toTable("my_catalog.my_schema.user_events")
)
同时,设置自动压缩属性,让 Databricks 在写入时自动合并:
ALTER TABLE my_catalog.my_schema.user_events
SET TBLPROPERTIES (delta.autoOptimize.autoCompact = true);
问题三:集群启动慢与自动缩放失灵
有段时间团队抱怨:每次打开 Notebook 都要等 5 分钟集群才能就绪。而且自动缩放(Auto-scaling)看起来是开着的,但集群规模从来没缩下来过,成本账单居高不下。
我排查后发现两个原因:
- 我们用的是按需实例(On-Demand),没有配置实例池(Instance Pool),每次扩容都要从零启动虚拟机。
- 自动缩放的最小节点数设成了 5,最大设成了 20,但日常负载只需要 2 个节点。
我的修复方案:
首先,创建实例池。在 Databricks 工作区菜单中进入 Compute -> Instance Pools,创建一个池:
- 池名称:
etl-worker-pool - 节点类型:
Standard_E4ds_v4(根据你的工作负载选择) - 最小空闲实例:3
- 最大容量:50
然后在集群配置中引用这个池:
{
"instance_pool_name": "etl-worker-pool",
"autoscale": {
"min_workers": 2,
"max_workers": 15
},
"spark_version": "14.3.x-scala2.12"
}
对于交互式开发场景,我改用 Serverless SQL Warehouse。它的冷启动在 10 秒以内,按查询计费,不用时自动归零:
-- 在 SQL Editor 中直接切换到 Serverless Warehouse
SELECT * FROM my_catalog.my_schema.sales_transactions LIMIT 10;
成本对比: 优化前我们每月的集群费用约 4200 美元。配置实例池 + 调整自动缩放参数后,降到了 1800 美元。省下来的钱够再跑三个实验集群。
问题四:Unity Catalog 权限报错
迁移到 Unity Catalog 后,整个团队被权限问题折磨了一周。最常见的报错是:
PermissionDeniedException: USER does not have SELECT on Table
或者更隐蔽的:
Catalog 'my_catalog' does not exist
明明在 UI 上能看到这个 Catalog,但代码里就是访问不了。
根本原因: Unity Catalog 的权限模型是三层继承的——Catalog → Schema → Table。只授予表级权限是不够的,用户还需要对上层 Catalog 和 Schema 有 USE CATALOG 和 USE SCHEMA 权限。
正确的授权脚本:
-- 第一步:授予 Catalog 访问权
GRANT USE CATALOG ON CATALOG my_catalog TO `data-team@company.com`;
-- 第二步:授予 Schema 访问权
GRANT USE SCHEMA ON SCHEMA my_catalog.my_schema TO `data-team@company.com`;
-- 第三步:授予具体表的查询权
GRANT SELECT ON TABLE my_catalog.my_schema.sales_transactions TO `data-team@company.com`;
如果你想让团队拥有某个 Schema 下所有表的权限,用 ALL TABLES:
GRANT SELECT ON ALL TABLES IN SCHEMA my_catalog.my_schema TO `data-team@company.com`;
一个容易忽略的坑: 如果你的存储路径在 ADLS Gen2 或 S3 上,还需要配置外部位置(External Location)的权限。否则即使有表权限,底层读不到文件还是会报错:
GRANT READ FILES ON EXTERNAL LOCATION 'abfss://data@storage.dfs.core.windows.net/raw/'
TO `data-team@company.com`;
问题五:Spark 作业 OOM 与数据倾斜
回到我凌晨那次事故的另一个原因:数据倾斜。我们的 sales_transactions 表按 region_code 分区,其中 region_code = 'CN' 的数据占了总量的 72%。当执行 GROUP BY region_code 时,处理 'CN' 的那个 Task 内存直接爆了。
诊断方法: 在 Spark UI 的 Stages 页面查看 Task 分布。如果某个 Task 处理的数据量远大于其他 Task,就是倾斜。
我的解决方案:
对于临时查询,用加盐(salting)技术打散热点键:
from pyspark.sql import functions as F
# 添加随机盐值(0-9)
df_salted = df.withColumn("salted_key",
F.concat(F.col("region_code"),
F.lit("_"),
F.floor(F.rand() * 10)))
# 按加盐键分组聚合
result = df_salted.groupBy("salted_key").agg(F.sum("amount").alias("total"))
# 去掉盐值,二次聚合
final = result.withColumn("region_code", F.split(F.col("salted_key"), "_")[0]) \
.groupBy("region_code").agg(F.sum("total").alias("total_amount"))
对于持久化表,重新选择分区键。我最终把分区键从 region_code 改成了 event_date,因为日期的分布更均匀:
-- 创建新表
CREATE TABLE my_catalog.my_schema.sales_transactions_v2
USING DELTA
PARTITIONED BY (event_date)
AS SELECT * FROM my_catalog.my_schema.sales_transactions;
-- 原子替换
ALTER TABLE my_catalog.my_schema.sales_transactions
RENAME TO sales_transactions_old;
ALTER TABLE my_catalog.my_schema.sales_transactions_v2
RENAME TO sales_transactions;
同时,增加倾斜分区的 shuffle 分区数:
spark.conf.set("spark.sql.shuffle.partitions", "800")
spark.conf.set("spark.sql.adaptive.skewedJoin.enabled", true) # AQE 自动处理倾斜
Databricks Runtime 14.x 默认开启了自适应查询执行(AQE),它能在运行时自动检测并处理倾斜 Join。但前提是你得让 shuffle.partitions 足够大,否则 AQE 没有足够的粒度来拆分。
实战总结与诚实评估
经过这些踩坑,我现在的 Databricks 项目启动清单是这样的:
| 检查项 | 默认配置 | 推荐配置 |
|---|---|---|
| shuffle.partitions | 200 | 800+(大表场景) |
| delta.autoOptimize | 未启用 | true |
| delta.logRetentionDuration | 30 days | 7 days(高频写入表) |
| 集群自动缩放最小节点 | 5 | 2 |
| 实例池 | 未使用 | 必须配置 |
| Unity Catalog 权限 | 仅表级 | Catalog + Schema + 表级 |
Databricks 的局限性我也必须说清楚:
- 冷启动问题无法完全消除。 即使配置了实例池,首次启动仍需 1-2 分钟。Serverless 模式可以缓解,但目前只支持 SQL Warehouse,不支持通用 Spark 集群。
- Unity Catalog 的权限模型学习曲线陡峭。 从传统的 IAM 角色迁移到 Unity Catalog 的三层模型,需要重新设计权限架构,这不是一个周末能搞定的事。
- 跨云协作仍然复杂。 虽然像 AWS SageMaker 和 Databricks Unity Catalog 的集成方案已经出现(通过 OpenSharing API),但配置跨平台的网络连通性和身份认证仍然需要大量基础设施工作。
- 成本控制需要持续关注。 Databricks 很容易用得爽但花得快。我强烈建议设置预算告警,并在非生产环境使用 Job Cluster 而非 All-Purpose Cluster——前者的单价大约是后者的一半。
凌晨那次事故后,我花了两天时间把上面这些配置全部落地。现在同样的 ETL 作业,运行时间从 3 小时降到了 12 分钟,集群成本降了 60%,而且再也没有因为并发写入或数据倾斜在半夜被叫醒过。
工具本身不会替你思考,但理解了它的脾气,你至少能睡个好觉。