Pandas转Spark DataFrame过慢,如何强制其驻留Driver节点?
问题描述
我们在Spark集群的Worker节点运行CPU密集型工作负载,通过RDD .collect()把计算结果拉到内存配置更高的Driver节点,后续处理生成Pandas DataFrame(这部分逻辑没法改)。现在需要把这个Pandas DataFrame存到Databricks,当前做法是转成Spark DataFrame后调用.saveAsTable()。
遇到的问题是:针对900列、5万行的表,Pandas转Spark DataFrame要花5分钟,而Spark DataFrame写入Databricks只需要10秒。推测慢的原因是转换时数据被自动分发到了集群Worker节点。
已经试过的无效操作:
- 设置
spark.default.parallelism=1,但日志显示数据还是被发到了Worker节点; - 用
repartition()或coalesce(),但这俩得在DataFrame创建后用,解决不了转换阶段的耗时。
想知道有没有办法强制Spark DataFrame只在Driver节点驻留、不做分布式处理,同时还能用到它方便的写入API?
当前写入代码:
delta_frame.write \ .mode("append") \ .option("delta.columnMapping.mode", "name") \ .option("mergeSchema", "true" if merge_schema else "false") \ .option("path", target_path) \ .partitionBy(partition_cols) \ .saveAsTable(full_table_name)
解决方案
方法1:创建时指定单分区+开启Arrow加速
Pandas转Spark DataFrame的耗时,很大程度来自数据序列化和分区分发。按下面两步优化能直接解决问题:
- 开启PySpark Arrow加速:在构建Spark Session时加上这个配置,大幅减少Pandas和Spark之间的数据转换开销:
spark = SparkSession.builder \ .appName("DriverOnlyDF") \ .config("spark.sql.execution.arrow.pyspark.enabled", "true") \ .getOrCreate() - 强制单分区创建DataFrame:调用
spark.createDataFrame()时传入numPartitions=1,让Spark直接在Driver节点完成转换,数据不会被分发到Worker:
这个操作从根源上避免了转换阶段的数据分发,能把转换耗时压下来。delta_frame = spark.createDataFrame(pandas_df, numPartitions=1)
方法2:本地文件中转(极端场景备选)
如果方法1还是没挡住数据分发,可以先把Pandas DataFrame写到Driver本地的临时文件(比如Parquet),再用Spark读取这个本地文件生成单分区DataFrame:
# 写到Driver本地临时路径 temp_path = "/tmp/local_temp_data.parquet" pandas_df.to_parquet(temp_path) # 读取本地文件生成Spark DataFrame(默认单分区,数据在Driver) delta_frame = spark.read.parquet(f"file://{temp_path}")
之后正常执行你的写入代码就行。
注意点
- 写入时用了
partitionBy的话,Spark会按分区列重分区,但这个过程的开销比转换阶段的分发小很多; - 确保Driver节点内存够装下整个Pandas DataFrame,别出现内存溢出。
内容的提问来源于stack exchange,提问作者DaManJ
相关产品推荐
相关产品推荐

