【问题标题】:Efficient execution on PySpark/Delta dataframes在 PySpark/Delta 数据帧上高效执行
【发布时间】:2019-11-01 13:22:06
【问题描述】:

在 Databricks 上使用 pyspark/Delta 湖,我有以下场景:

sdf = spark.read.format("delta").table("...")
result = sdf.filter(...).groupBy(...).agg(...)

analysis_1 = result.groupBy(...).count() # transformation performed here
analysis_2 = result.groupBy(...).count() # transformation performed here

据我了解,使用 Delta 湖的 Spark,由于链式执行,result 实际上不是在声明时计算,而是在使用时计算。

然而,在这个例子中,它被多次使用,因此最昂贵的转换被多次计算。

是否可以在代码中的某个点强制执行,例如

sdf = spark.read.format("delta").table("...")
result = sdf.filter(...).groupBy(...).agg(...)
result.force() # transformation performed here??

analysis_1 = result.groupBy(...).count() # quick smaller transformation??
analysis_2 = result.groupBy(...).count() # quick smaller transformation??

【问题讨论】:

标签: apache-spark-sql databricks delta-lake


【解决方案1】:

在我看来,问题无处不在,或者说不清楚。但是,如果您是 Spark 的新手,则可能会出现这种情况。

所以:

关于 .force 的使用,请参阅https://blog.knoldus.com/getting-lazy-with-scala/ .force 将不适用于数据集或数据框。

这与 pyspark 或 Delta Lake 方法有关吗?不,不。

analysis_1 = result.groupBy(...).count() # quick smaller transformation?? 
  • 这实际上是一个在转换之前最有可能导致洗牌的动作。

所以,我认为您的意思是正如我们尊敬的 pault 所说的那样:

  • .cache 或 .persist

我怀疑你需要:

result.cache 

这意味着您的第二次行动分析_2不需要重新计算一路回到此处显示的源代码

(2) Spark Jobs
Job 16 View(Stages: 3/3)
Stage 43: 
8/8
succeeded / total tasks 
Stage 44: 
200/200
succeeded / total tasks 
Stage 45:   
1/1
succeeded / total tasks 
Job 17 View(Stages: 2/2, 1 skipped)
Stage 46: 
0/8
succeeded / total tasks skipped
Stage 47: 
200/200
succeeded / total tasks 
Stage 48:   
1/1
succeeded / total tasks 

随着对 Spark 的改进,shuffle 分区被保留,在某些情况下也会导致跳过阶段,尤其是对于 RDD。对于数据帧,需要缓存才能获得我观察到的跳过阶段效果。

【讨论】:

  • 我确实是 Spark 的新手,所以我试图掌握各种概念。但是,它正在取得进展:-)谢谢您的解释。正如@pault 暗示的那样,我正在寻找的是cache 的概念。谢谢大家。
猜你喜欢
  • 1970-01-01
  • 2021-10-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-10-21
  • 1970-01-01
相关资源
最近更新 更多