【问题标题】:Spark cache function: Caching Job and Caching StageSpark缓存功能:缓存作业和缓存阶段
【发布时间】:2018-08-22 01:50:54
【问题描述】:

我是 spark 新手,可以在这里使用一些指导。 我们有一些基本代码要读入 csv、缓存并输出到 parquet:

   1. val df=sparkSession.read.options(options).schema(schema).csv(path)
   2. val dfCached = df.withColumn()....orderBy(some Col).cache()
   3. dfCached.write.partitionBy(partitioning).parquet(outputPath)

AFAIK,一旦我们调用 parquet 调用(一个动作),应该执行缓存命令以在应用该动作之前保存 DF 的状态。

在我看到的 spark UI 中:

  1. 执行上述 #2 中的 cache 调用的单个分阶段作业
  2. 然后是一个正在执行parquet 调用的作业。这项工作有2个阶段; 1 似乎在重复缓存步骤,而第二个执行转换为镶木地板。 (见下图)

为什么我有缓存 Job 和缓存 Stage? 我希望只有一个或另一个,但似乎我们在这里缓存了两次。

【问题讨论】:

  • 你在哪个 spark 版本上运行这个?
  • @sai Spark 版本 2.3.0

标签: apache-spark


【解决方案1】:

我不是 100% 确定,但似乎正在发生以下情况:

  1. 加载 csv 数据时,它将在工作节点之间拆分。我们调用 cache() 并且每个节点将它接收到的数据存储在内存中。这是第一个缓存作业。

  2. 当我们调用 partitionBy(...)data 时,需要根据传递给函数的 args 在不同的执行器之间重新组合。由于我们正在缓存数据并且数据已经从一个执行器移动到另一个执行器,因此我们需要重新缓存已洗牌的数据。这得到了证实,因为第二个缓存阶段显示了一些随机写入数据。此外,缓存阶段显示的任务比初始缓存作业少;可能是因为只有混洗后的数据需要重新缓存,而不是整个数据帧。

  3. 调用 parquet 阶段。我们可以看到一些 shuffle 读取数据,这些数据显示 executor 正在读取新的 shuffle 数据。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-05
    • 2018-02-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多