PySpark中collect()预热比show()更能加速查询的原因探究
问题描述
近期发现,在Spark中对转换操作执行collect()进行预热后,后续查询的执行速度会显著加快;但使用show()进行预热时,却对查询性能没有影响。附上复现代码,其中run_collect_to_warm_up函数的耗时明显低于另一个函数,求解该现象的原因。
复现代码:
from time import perf_counter import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * spark = SparkSession.builder.appName("TEST").getOrCreate() sdf = spark.range(0, 1000000).withColumn( 'id', col('id') ).withColumn('v', rand()) @pandas_udf(DoubleType()) def pandas_plus_one(pdf): return pdf + 1 def run_collect_to_warm_up(): res = sdf.select(pandas_plus_one(col("v"))) res.collect() st = perf_counter() for _ in range(10): res.show() print(f"Time (warming up with collect): {(perf_counter() - st) * 1000}") def run_show_to_warm_up(): res = sdf.select(pandas_plus_one(col("v"))) res.show() st = perf_counter() for _ in range(10): res.show() print(f"Time (warming up with show): {(perf_counter() - st) * 1000}") run_show_to_warm_up() run_collect_to_warm_up()
原因解析
这个差异核心在于collect()和show()的执行行为完全不同,结合Spark对Pandas UDF的运行逻辑,具体原因如下:
数据处理范围差异
show()默认仅处理并返回前20条数据,只会触发1个(或极少数)分区的计算,Pandas UDF只被调用一次处理这小批量数据;而collect()会拉取全量数据,触发所有分区的计算,让Pandas UDF在所有分区上完成首次执行。JIT编译与初始化工作的覆盖范围
Pandas UDF依赖PyArrow和Pandas的向量化操作,首次运行时需要完成:- Python与JVM间序列化/反序列化逻辑的初始化
- Pandas/NumPy向量化操作的JIT编译
- Executor侧相关类加载、资源分配
collect()的全量执行会把这些初始化工作在预热阶段全部完成,后续show()直接复用已就绪的资源和编译后的逻辑;而show()预热只完成了极小范围的初始化,后续查询仍要为剩余分区重复执行这些耗时操作。
执行计划与元数据的缓存
collect()触发全量计算后,Spark会完整缓存该查询的执行计划、分区元数据等信息,后续相同查询的show()可以直接复用,省去重复的计划解析、优化步骤;而show()的小数据量执行可能不会触发完整的元数据缓存,后续查询仍需重新解析计划。Executor资源激活程度不同
collect()会让所有Executor参与计算,提前激活CPU、内存资源,建立稳定的数据传输连接;而show()可能仅用到单个Executor的少量资源,后续查询时其他Executor仍需启动资源、建立连接,额外增加延迟。
内容的提问来源于stack exchange,提问作者jiashenC
相关产品推荐
相关产品推荐

