优化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
相关产品推荐
相关产品推荐

