【发布时间】: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