Databricks中DataFrame保存过慢求助:数千行数据耗时超1小时
优化Databricks小数据集写入速度的建议
核心问题在于你使用了coalesce(1),这会强制将所有数据合并到单个分区,完全抵消了Spark的并行处理优势,单线程写入几千行数据也会因为Spark作业的调度开销变得异常缓慢,更不用说重复执行数千次的场景。以下是具体优化方案:
1. 移除coalesce(1),利用并行写入
Spark会根据数据量自动分配合理的分区数,几千行数据的默认分区设置足以支撑快速并行写入。如果后续工具要求单个CSV文件,可以在写入完成后用dbutils.fs合并文件,而非在写入阶段强制单分区。
修改后的写入代码:
df1.write.format("com.databricks.spark.csv")\ .option("header", "true")\ .save("dbfs:/FileStore/user/me/logpos.csv")
2. 减少不必要的数据读取与处理
原代码中先用select *读取全量字段,再通过select筛选所需列,额外增加了数据传输和内存开销。直接在SQL查询中只选择需要的字段,一步到位得到目标DataFrame:
table_name = 'coordinates' df = spark.sql(f""" select character.accountid, character.health from main_frame where event_date = '2023-01-18' and session_unique_id = '34eb1a29-aebb-4425-9997-7074e45244e9' """)
3. 批量处理替代单次循环执行
由于需要重复执行数千次,避免逐个处理单个session或日期的小任务。可以一次性读取多个目标条件的数据,按session_unique_id或event_date分区后批量写入,减少Spark作业的重复调度开销。例如:
# 假设需要处理多个session ID,放在列表中 session_ids = ["id1", "id2", "..."] df = spark.sql(f""" select character.accountid, character.health, session_unique_id from main_frame where event_date = '2023-01-18' and session_unique_id in ({','.join([f"'{id}'" for id in session_ids])}) """) # 按session_unique_id分区写入,每个session对应一个目录 df.write.format("com.databricks.spark.csv")\ .option("header", "true")\ .partitionBy("session_unique_id")\ .save("dbfs:/FileStore/user/me/logpos/")
4. 切换到更高效的存储格式
如果后续流程不强制要求CSV,建议使用Parquet或Delta格式存储。这类列式存储的写入、读取效率远高于CSV,且支持Schema校验、增量更新等特性,适合机器学习数据的反复调研与处理:
df.write.format("delta")\ .save("dbfs:/FileStore/user/me/logpos_delta")
5. 优化Spark资源配置
针对批量执行场景,适当调整Spark作业的资源参数,比如增加executor数量(无需过大,几千行数据2-4个executor足够),减少作业调度等待时间。可以在 notebook 开头设置:
spark.conf.set("spark.executor.instances", "2") spark.conf.set("spark.executor.cores", "2")
内容的提问来源于stack exchange,提问作者albusdemens
相关产品推荐
相关产品推荐

