PySpark性能问题:缓存临时视图后仍重复查询原始数据集
PySpark性能问题:缓存临时视图后仍重复查询原始数据集
我来帮你分析下问题出在哪,以及怎么快速解决~
问题根源:缓存操作根本没生效
你代码里的这个写法藏着一个关键坑:
df_data.cache().createOrReplaceTempView("subscriptionHistory")
PySpark的DataFrame是不可变对象,cache()方法会返回一个带缓存标记的新DataFrame,但你没把这个新对象重新赋值给df_data。所以原来的df_data还是没有缓存标记的版本,后续的display(df_data)和groupBy操作都不会用缓存,而是重新执行从读Parquet开始的全量计算——这就是为啥第二次操作还是慢到离谱。
另外补充下:createOrReplaceTempView()返回的是None,哪怕你写成df_data = df_data.cache().createOrReplaceTempView(...),也会把df_data变成空值,完全行不通。
正确的优化步骤
1. 正确触发缓存
先把带缓存标记的DataFrame重新赋值,再执行action操作(比如display())触发缓存落地:
# 给df_data加上缓存标记并重新赋值 df_data = df_data.cache() # 创建临时视图 df_data.createOrReplaceTempView("subscriptionHistory") # 执行action触发缓存(display属于action操作) display(df_data)
这一步跑完,你的93条数据就会被存在Spark的存储层里,后续操作直接读缓存就行。
2. 用缓存数据执行分组操作
你有两种方式利用缓存好的数据:
- 方式一:直接用已缓存的
df_data做分组:
import pyspark.sql.functions as f df_result = df_data.groupBy("subscriptionState").agg(f.min("dateTime"), f.max("dateTime")) display(df_result)
- 方式二:通过临时视图写SQL查询(有时候SQL优化器会更智能):
df_result = spark.sql(""" SELECT subscriptionState, MIN(dateTime) AS start_time, MAX(dateTime) AS end_time FROM subscriptionHistory GROUP BY subscriptionState """) display(df_result)
3. 额外的性能提速小技巧
除了缓存,这些方法也能帮你进一步优化:
- 检查Parquet分区:如果S3上的Parquet数据是按
supplyUuid或dateTime分区的,直接用分区过滤能跳过全表扫描,比如:# 如果按supplyUuid分区,Spark会自动跳过无关分区,速度会快很多 df_data = spark.read.parquet(baseName).filter("supplyUuid='xxxxxxxxx...'") - 调整缓存存储级别:集群内存够的话用默认的
MEMORY_ONLY就行;内存紧张可以用MEMORY_AND_DISK,避免缓存被轻易清除:from pyspark.storagelevel import StorageLevel df_data = df_data.persist(StorageLevel.MEMORY_AND_DISK) - 用完缓存记得释放:后续不需要这个数据时,调用
df_data.unpersist()释放集群内存,避免资源浪费。
验证缓存是否生效
你可以通过Spark UI确认:
- 打开集群的
4040端口(默认Spark UI端口) - 进入Storage页面,查看是否有
subscriptionHistory的缓存条目,以及缓存数据量是否是93条。
这样操作后,你的分组查询应该就能秒出结果啦~
备注:内容来源于stack exchange,提问作者deechean wang
相关产品推荐
相关产品推荐

