Databricks写入小型DataFrame耗时过长问题求助
Spark写入Delta格式耗时过长问题分析
问题场景
- 源数据:DB2表,900万行、40列,执行
count操作仅需50秒 - 操作:每月全量加载数据,按
execution_date分区写入Delta表,输出文件总大小约600MB,但写入耗时数小时 - 执行代码:
spark.conf.set("spark.sql.sources.partitionOverwriteMode","dynamic") SCRIPT = "select * from DB.TABLE" DF = spark.read.format("jdbc").option("url", url).option("driver", props["driver"]).option("user", props["username"]).option("password", props["password"]).option("query", SCRIPT).load() DF = DF.withColumn("DE_INSERTED_DATE", current_date()).withColumn("execution_date", lit(last_day_of_month)) DF.write.format("delta") \ .option("path", "s3://PATH-TO-MY-TABLE") \ .mode("overwrite") \ .partitionBy("execution_date") \ .saveAsTable("DATABASE.TABLE") #insertInto("DATABASE.TABLE")
- 补充信息:集群详情及监控截图见附件
可能的问题点及优化方案
1. JDBC读取并行度不足
虽然count操作快,但JDBC默认可能单线程读取DB2数据,导致后续写入阶段数据处理并行度不够,拖慢整体流程。
- 优化方式:给JDBC读取添加分区参数,利用Spark多并行度读取数据:
DF = spark.read.format("jdbc") .option("url", url) .option("driver", props["driver"]) .option("user", props["username"]) .option("password", props["password"]) .option("query", SCRIPT) .option("numPartitions", 20) # 根据集群核数调整,建议10-30之间 .option("partitionColumn", "数值型主键/日期列") # 选一个分布均匀的列 .option("lowerBound", "列最小值") .option("upperBound", "列最大值") .load()
2. Delta写入的文件大小与并行度不匹配
600MB数据若生成大量小文件,会导致S3写入时IO开销剧增。默认的spark.sql.shuffle.partitions(200)会生成过多小文件。
- 优化方式:
- 调整shuffle分区数:
spark.conf.set("spark.sql.shuffle.partitions", 50),让每个分区文件大小约12MB,符合Delta的合理文件范围 - 写入前手动调整数据分区:
DF.repartition(50).write...,减少小文件生成 - 启用Delta自动优化:添加
.option("optimizeWrite", "true"),让Delta自动合并小文件
- 调整shuffle分区数:
3. S3存储IO瓶颈
跨区域写入、S3连接数不足或未启用快速上传都会导致写入延迟。
- 优化方式:
- 确保集群区域与S3存储区域一致,避免跨区域网络延迟
- 调整S3相关配置:
spark.conf.set("spark.hadoop.fs.s3a.connection.maximum", 100) spark.conf.set("spark.hadoop.fs.s3a.fast.upload", "true")
4. 写入模式与分区策略冲突
使用mode("overwrite")配合saveAsTable,即使开启了动态分区覆盖,仍可能触发全表元数据操作或扫描,增加不必要的开销。
- 优化方式:改用
insertInto(取消代码中注释),确保仅覆盖目标execution_date分区,避免全表操作:
DF.write.format("delta") \ .option("path", "s3://PATH-TO-MY-TABLE") \ .mode("overwrite") \ .partitionBy("execution_date") \ .insertInto("DATABASE.TABLE")
注意:使用insertInto需保证DataFrame列与目标表完全匹配。
5. 集群资源瓶颈
结合监控截图确认:
- 若CPU持续满载:增加Executor核数或调整任务并行度
- 若频繁GC:调大
spark.executor.memory和spark.driver.memory - 若磁盘IO过高:检查本地磁盘是否存在瓶颈,或调整存储相关配置
内容的提问来源于stack exchange,提问作者finman
相关产品推荐
相关产品推荐

