You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

优化PySpark DataFrame写入ADLS Delta表耗时的方法咨询

问题描述

我在Azure Databricks上使用**1个驱动节点、1-8个工作节点(每个节点配备16核CPU与56GB内存)**的多节点集群运行Notebook,从Azure ADLS读取3万条源数据。Notebook包含若干数据转换步骤及两个业务必需的UDF,所有转换步骤仅耗时12分钟(符合预期),但将最终DataFrame写入ADLS Delta表却耗时超2小时。

相关代码片段

转换完成后的显示与持久化

# 所有数据读取和转换代码在此之前
# 仅在保存到delta表前有一个display语句,到这一步仅耗时12分钟
data.display() 
# 持久化DataFrame
from pyspark import StorageLevel
data.persist(StorageLevel.MEMORY_ONLY)

写入Delta表的核心代码(耗时超2小时)

# 持久化品牌提取结果
(
  data
  .write
  .format('delta')
  .mode('overwrite')
  .option('overwriteSchema', 'true')
  .saveAsTable('output_table')
)

另一种尝试的写入方式(无明显改善)

mount_path = "\/mnt\/********\/"
table_name = "********"
adls_path = mount_path + table_name
(data.write.format('delta').mode('overwrite').option('overwriteSchema', 'true').save(adls_path))
优化写入耗时的方法
  • 调整数据分区:先执行data.rdd.getNumPartitions()查看当前分区数,若分区数远大于集群可用核心数(比如超过64),用data.repartition(8)或data.coalesce(8)合并分区(8对应工作节点数,可根据实际集群配置调整),避免生成大量小文件拖慢写入。
  • 优化持久化策略:当前使用MEMORY_ONLY,若数据无法完全存入内存会溢写到磁盘,反而降低效率。可改为MEMORY_AND_DISK,或者直接移除持久化代码——3万条数据量极小,写入时重新计算可能比从磁盘读取更快。
  • 关闭不必要的Schema选项:仅当Schema确实需要更新时保留.option('overwriteSchema', 'true'),若每次写入Schema无变化,移除该选项,减少元数据校验和操作的开销。
  • 排查存储层问题:确认ADLS挂载点权限配置正确,无网络瓶颈。可尝试将数据写入DBFS本地路径测试速度,排查是否是ADLS端的性能问题。
  • 检查数据均衡性:执行data.groupBy(spark_partition_id()).count().show()查看各分区数据量是否均匀,若存在数据倾斜,对倾斜字段进行预处理(如加盐)后再写入。
  • 启用Delta优化选项:若写入的是静态数据,添加.option("dataChange", "false")减少Delta的元数据更新操作;写入完成后可执行OPTIMIZE output_table ZORDER BY <核心字段>(数据量小的情况下收益有限,按需使用)。

内容的提问来源于stack exchange,提问作者user22

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 20:20:30