关于PySpark DataFrame函数参数传递机制及repartition、persist性能作用的技术问询
问题
- 当将变量
df传入func2时,该参数是按引用传递还是按值传递? - 在func1中执行的
repartition("col1", "col2")操作,是否能够提升func2中groupBy操作的性能? - 若在func1中调用
persist(StorageLevel.MEMORY_AND_DISK),该DataFrame会被存储至内存(及磁盘)中吗?后续在func2中引用该df时,是否会直接读取func1中首次存储的缓存数据?
解答
1. 参数传递方式
在Spark里,DataFrame是不可变对象,你传入func2的其实是这个DataFrame对象的引用,但因为DataFrame的设计是不可修改的——所有对它的操作(比如select、groupBy)都会返回一个全新的DataFrame——所以实际表现上和「按值传递」的效果差不多。
举个例子:func2里对传入的df做select、groupby这些操作,都是生成新的DataFrame实例,完全不会影响func1里原来的那个df对象。所以不用担心在func2里操作会修改原数据。
2. repartition对groupBy性能的影响
这得分情况具体分析:
func2里的groupBy是针对col1字段,但func1里的repartition("col1", "col2")是按col1和col2两个字段的组合哈希值来分区的。这意味着,即使两条记录的col1相同,但col2不同,它们的哈希值可能不一样,会被分到不同的分区中。
当func2执行groupBy("col1")时,还是需要把分散在不同分区的同col1数据shuffle到同一个分区里才能聚合,所以这种repartition方式对这个groupBy的性能提升非常有限。
如果想优化这个groupBy的性能,更合理的做法是直接在func1里按col1单独repartition:df.repartition("col1")。这样同一个col1的记录会被分到同一个分区,后续groupBy的时候就不需要额外的shuffle操作,性能会明显提升。
3. persist的缓存逻辑与后续读取
首先要记住Spark的核心特性:惰性求值。persist(StorageLevel.MEMORY_AND_DISK)只是给DataFrame打上一个「需要缓存」的标记,并不会立刻把数据存入内存或磁盘——只有当这个DataFrame触发action操作(比如count()、show()、write()等)时,才会实际计算数据并执行缓存操作。
再看你的func1代码:你给原始的df调用了persist,但返回的是df.repartition("col1", "col2")生成的新DataFrame。这里有两个关键细节:
- 原始
df的缓存只会在它被首次计算(触发action)的时候才会被存储。如果func1里没有对原始df执行任何action,那缓存标记不会实际生效。 - 你返回的是repartition后的新DataFrame,这个新df并没有被标记persist。所以当func2里对这个新df进行操作时,会触发整个计算链路:原始df -> repartition -> func2的操作。这时候如果原始df还没被缓存,就会先计算原始df,然后缓存它(因为之前打了标记),再执行repartition;如果后续还有对这个链路的action操作,原始df会直接读取缓存,但repartition的步骤还是会重新执行(因为新df没被缓存)。
如果想让func2直接读取缓存后的分区数据,应该调整代码顺序,把persist移到repartition之后:
def func1(): # some code to create a dataframe df df = df.repartition("col1", "col2") df.persist(StorageLevel.MEMORY_AND_DISK) return df
这样返回的df就是被标记缓存的,当首次触发action时,会把repartition后的结果缓存起来,后续func2操作时就能直接读取缓存的数据了。
内容的提问来源于stack exchange,提问作者user3868051

