You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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数量的理解是否正确?

你的理解有部分偏差,核心纠正和确认如下:

  1. 原代码无action的场景:
    get_and_prepare_data()中的df.unpersist()会在函数返回前立即执行,但因为df未被缓存,所以无任何操作。当join_and_filter()调用count()时,触发的执行计划是:

    • Spark SQL查询 → repartition → groupBy/agg → filter → count
      这里unpersist不会被纳入DAG执行步骤,因为它不是数据转换操作,只是缓存管理命令,和数据计算流程无关。
  2. 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.

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 18:18:38