Spark中未缓存DataFrame调用unpersist的作用及执行流程问询
关于Spark DataFrame unpersist与懒执行的问题
代码背景
数据处理函数
def get_and_prepare_data(): # Read the data input_df = spark.sql('SELECT * FROM my_table') try: # Do a bunch of transformations, etc df = input_df.repartition(partition_ct) df = df.groupBy(order_id).agg(collect_list(item_id).alias("my_items")) except Exception as err: error_list.append(str(err)) log_error(err) finally: df.unpersist() return df
后续调用函数
def join_and_filter(df): filtered = df.filter(...) filtered.count()
问题1:调用df.unpersist()是否有实际作用?是否存在副作用?
- 对从未调用过
cache()或persist()的DataFrame执行unpersist()没有任何实际作用,因为Spark的缓存系统中根本没有这个DataFrame的缓存条目,自然没有内容可移除。 - 这种调用也不存在副作用:Spark内部会直接忽略该操作,不会抛出异常,也不会影响后续对该DataFrame的任何转换或action操作。
问题2:关于懒执行与Job数量的理解是否正确?
你的理解有部分偏差,核心纠正和确认如下:
原代码无action的场景:
get_and_prepare_data()中的df.unpersist()会在函数返回前立即执行,但因为df未被缓存,所以无任何操作。当join_and_filter()调用count()时,触发的执行计划是:- Spark SQL查询 → repartition → groupBy/agg → filter → count
这里unpersist不会被纳入DAG执行步骤,因为它不是数据转换操作,只是缓存管理命令,和数据计算流程无关。
- Spark SQL查询 → repartition → groupBy/agg → filter → count
group by后添加
count()的场景:- 第一个
count()会触发JOB A:SQL查询→repartition→groupBy/agg→count,这是完整的计算Job。 - 之后的
unpersist()依然无作用(因为df没被缓存)。 - 当
join_and_filter()调用count()时,会触发JOB B:SQL查询→repartition→groupBy/agg→filter→count,确实是第二个独立Job。 - 只有在groupBy/agg后添加
df.cache()(或persist()),JOB A计算后的结果才会被缓存,JOB B会直接读取缓存数据,避免重复计算。
- 第一个
所以你关于“无缓存时两次action会触发两个独立Job”的结论是正确的,但错误地把unpersist纳入了DAG执行步骤,这点需要纠正。
内容的提问来源于stack exchange,提问作者S.S.
相关产品推荐
相关产品推荐

