Spark DataFrame未缓存场景下执行相关问题咨询
嘿,咱们先把核心概念掰明白:PySpark是惰性求值的——也就是说,你写的那些DataFrame转换操作(比如spark.sql("select * from old_df")、withColumn)并不会立刻执行计算,只是在构建一个「执行计划(lineage)」,只有当你触发**动作(Action)**操作(比如count()、show())的时候,Spark才会真正跑起来计算数据。
问题1:执行count()命令后,new_df是否实际存在?
首先,new_df = spark.sql("select * from old_df")这行代码执行完,new_df就已经「存在」了——但它不是存在内存/磁盘里的实际数据集,而是一个包含了查询逻辑的执行计划。
当你执行print(new_df.count())时,count()是一个动作操作,会触发Spark去执行这个查询,计算出结果并返回行数,但执行完之后,Spark不会把new_df对应的数据集持久化下来。new_df本身还是那个执行计划,下次你再对它调用动作操作,Spark还是会重新跑一遍查询。
简单说:new_df作为逻辑上的DataFrame一直存在,但它对应的实际数据不会因为一次count()就被保存下来。
问题2:若将count()替换为show(5),问题1的答案是否会改变?
答案是不会。
show(5)同样是一个动作操作,它会触发Spark执行查询,取出前5条数据显示出来,但和count()一样,执行完之后也不会把new_df的数据集持久化。new_df依然只是那个执行计划,本质和之前没区别。
问题3:创建new_df的初始步骤是否会被重新执行?
会的!
当你执行new_df = new_df.withColumn('foo', 新列计算公式)时,你是在原来的new_df的执行计划基础上,添加了一个「新增列」的转换步骤,生成了一个新的执行计划(现在的new_df指向这个新计划)。
当你再次调用print(new_df.count())这个动作操作时,Spark会从头执行整个执行计划:先跑初始的select * from old_df,再执行withColumn的计算,最后统计行数。
如果不想让初始步骤重复执行,你可以用cache()或者persist()方法把中间结果持久化,比如:
new_df = spark.sql("select * from old_df") new_df.cache() # 把new_df的执行结果缓存起来 print(new_df.count()) # 第一次执行会跑查询并缓存 new_df = new_df.withColumn('foo', 新列计算公式) print(new_df.count()) # 这里初始的select步骤就不会重新跑了,直接用缓存的数据
内容的提问来源于stack exchange,提问作者B_Miner

