Databricks环境Spark作业写入分区Delta表卡住中止问题求助
故障排查步骤
- 查看Spark UI的Stage详情,定位卡住的任务阶段:如果是写入阶段出现长尾task,多数是数据倾斜或者分区粒度过细导致的小文件爆炸问题;如果是shuffle阶段卡住,检查shuffle分区配置是否合理。
- 检查作业运行日志,确认中止的报错原因:是executor OOM、driver OOM,还是资源调度超时。分区总数过高的场景下,driver需要管理的元数据量会大幅上升,很容易出现driver内存不足的问题。
- 统计最终的分区总数:按A(8419个唯一值)、year、month三个字段分区,假设年月有10个组合,总分区数就会达到8万+,远高于Delta表建议的最大分区数(建议不超过1万),过多的分区会导致写入时元数据更新耗时指数级增长,同时产生大量小于128MB的小文件,严重拖慢写入效率。
- 确认源DataFrame的分区数:如果读取CSV后的DataFrame分区数过少,单个task需要处理的数据量过大,也会导致写入速度过慢。
解决方案
- 优先调整分区策略:当前分区粒度过细,不符合Delta表的最佳实践(单分区数据量建议不低于1GB)。建议仅保留
year、month作为分区字段,将A字段设置为ZORDER索引,查询时对A字段过滤的性能和分区基本一致,同时可以将分区数降低到数十个级别,大幅降低写入开销。 - 若必须保留原分区规则,写入前做重分区:写入前执行
df = df.repartition("A", "year", "month"),保证同一分区的数据被分配到同一个task处理,避免大量task跨分区写产生海量小文件。也可以指定重分区的数量,1.8亿条数据建议设置400-800个重分区,单个task处理500MB-1GB数据,性能最优。 - 开启Delta表写入优化配置:写入前设置以下Spark参数:
# 开启自动优化写入,合并小文件 spark.conf.set("spark.databricks.delta.optimizeWrite.enabled", "true") # 写入后自动 compact 小文件 spark.conf.set("spark.databricks.delta.autoCompact.enabled", "true") # 调整shuffle分区数匹配数据量 spark.conf.set("spark.sql.shuffle.partitions", 600)
- 调整作业资源配置:给driver分配至少16GB内存(管理大量分区元数据需要更大内存),executor分配8-16GB内存、4-8核,保证资源足够支撑大量分区的写入操作。
- 全量写入场景可先写临时表:如果是首次全量导入数据,先使用
mode("overwrite")写入临时Delta表,验证数据无误后再通过append方式写入正式表,避免重复的元数据校验开销。
内容的提问来源于stack exchange,提问作者Abhishek Singh
相关产品推荐
相关产品推荐

