【问题标题】:Does skipped stages have any performance impact on Spark job?跳过的阶段对 Spark 作业有任何性能影响吗?
【发布时间】:2020-07-27 02:03:56
【问题描述】:

我正在运行一个 Spark 结构化流式传输作业,该作业涉及创建一个空数据帧,并使用每个微批处理对其进行更新,如下所示。每执行一次微批处理,阶段数就会增加 4。为了避免重新计算,我在循环内的每次更新后将更新的 StaticDF 持久化到内存中。这有助于跳过每个新微批次创建的额外阶段。

我的问题 -

1) 即使完成的阶段总数保持不变,因为增加的阶段总是被跳过,但它是否会导致性能问题,因为在一个时间点可能有数百万个跳过的阶段?
2)当缓存的RDD的一部分或全部不可用时会发生什么? (节点/执行器故障)。 Spark 文档说到目前为止它并没有具体化从多个微批次接收到的全部数据,这是否意味着它需要再次从 Kafka 读取所有事件以重新生成 staticDF?

// one time creation of empty static(not streaming) dataframe
val staticDF_schema = new StructType()
      .add("product_id", LongType)
      .add("created_at", LongType)
var staticDF = sparkSession
.createDataFrame(sparkSession.sparkContext.emptyRDD[Row], staticDF_schema)

// Note : streamingDF was created from Kafka source
    streamingDF.writeStream
      .trigger(Trigger.ProcessingTime(10000L))
      .foreachBatch {
        (micro_batch_DF: DataFrame) => {

        // fetching max created_at for each product_id in current micro-batch
          val staging_df = micro_batch_DF.groupBy("product_id")
            .agg(max("created").alias("created"))

          // Updating staticDF using current micro batch
          staticDF = staticDF.unionByName(staging_df)
          staticDF = staticDF
            .withColumn("rnk",
              row_number().over(Window.partitionBy("product_id").orderBy(desc("created_at")))
            ).filter("rnk = 1")
            .drop("rnk")
              .cache()

          }

【问题讨论】:

  • 这不就是检查点的用途吗? spark.apache.org/docs/latest/…
  • @mazaneicha 谢谢,我在这里寻找检查点文档,其中说它将偏移量和中间聚合存储到检查点位置。仍然不清楚它会将 staticDF 存储在 HDFS 上还是仅使用偏移量再次从 Kafka 读取所有内容。 spark.apache.org/docs/latest/…
  • @user10938362 不幸的是,我知道重新计算不会发生,因为结果已经可用,但我担心在创建 DAG 期间是否会有开销,因为总阶段数可能达到数百万(即使已完成阶段将始终保持 7 并且将跳过休息)
  • 检查点保存元数据和状态。

标签: scala apache-spark spark-streaming spark-structured-streaming spark-streaming-kafka


【解决方案1】:

即使跳过的阶段不需要任何计算,但我的工作在一定数量的批次后开始失败。这是因为每次批处理执行时 DAG 都会增长,使其无法管理并引发堆栈溢出异常。

为了避免这种情况,我不得不打破火花谱系,这样每次运行的阶段数都不会增加(即使它们被跳过)

【讨论】:

  • 如何打破火花血统?
  • 由于我的数据足够小,我将其存储在 scala 映射中,并使用 spark 上下文(而不是使用缓存的 rdd)在每个微批次中创建一个新的 RDD。这消除了对先前 RDD 的依赖来计算新的 RDD 并打破了沿袭。我还使用我在每个微批次中获得的值更新这个 scala 映射,以便它始终更新。
  • 嗨!我和你有同样的问题,我的工作有些跳过,我想问你如何处理这个问题?我还是没明白你说的 spark lineage 是什么意思,是不是解决了防止跳槽的问题?
  • @RudyTriSaputra 请发布一个单独的问题,详细说明您的问题。我的问题可能与你的完全不同。
  • 上帝感谢您回复我的评论,我有一个与性能问题有关的问题,而且我有一个跳过阶段的工作。这是我的问题link
猜你喜欢
  • 2014-05-15
  • 2016-04-23
  • 2011-11-12
  • 2014-01-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-11-04
  • 2015-07-14
相关资源
最近更新 更多