Spark函数内缓存DataFrame的设计合理性与内存管理问题咨询
Spark函数内缓存DataFrame的设计合理性与内存管理问题咨询
I have this toy example function returning a cached spark DataFrame (DF):
def foo(df): try: base = complicated_query(df) base.cache() # lazy cache base.count() # trigger cache - wrong design??? num1 = base.withColumn('number', f.lit('1')) num2 = base.withColumn('number', f.lit('2')) return num1.union(num2) finally: None # base.unpersist()
The purpose of foo is to simply encapsulate temporary variables (DataFrames) which I don't want to have in the outer scope. base is some complicated DF used twice, hence I cache it. My questions are around the explicit count call (suggested by ChatGPT). This is to trigger caching but I feel like it is a wrong design.
- Why do we need
count(an action) at this point? What is the gain? The actual caching would happen anyway on the first action to the return value offooif I didn't callcount. - I noticed that calling
footwice with the same inputdfhas different execution times: first is slow, which is expected since acountis called. But the second is almost instant. Why is that? Surelybaseis already cached and thencountis trivial, but thebasereference in the second run is a different reference than that in the first run. How does spark know that it can reuse a cached DF from the first run? (whose memory btw I leaked I guess since afterfooexits I cannotunpersistit.) - Do we really have memory leak when
fooexits? How can Iunpersistbase? Do I have to unpersist at all? I know there isspark.catalog.clearCache()to wipe out all cached DFs but I would like to do this explicitly forbase. That's why there is/was afinallyclause in the function, to prevent a leak, but that was a failed attempt as in that case I was freeing up the cache before I could even use it...
Can you please help resolve these?
备注:内容来源于stack exchange,提问作者Andras Vanyolos
相关产品推荐
相关产品推荐

