【问题标题】:Spark Structured Streaming checkpoint usage in productionSpark Structured Streaming 检查点在生产中的使用
【发布时间】:2020-10-29 04:05:26
【问题描述】:

在使用 Spark 结构化流时,我无法理解检查点的工作原理。

我有一个生成一些事件的 spark 进程,我将这些事件登录到 Hive 表中。 对于这些事件,我会在 kafka 流中收到确认事件。

我创建了一个新的火花过程

  • 将 Hive 日志表中的事件读取到 DataFrame 中
  • 使用 Spark Structured Streaming 将这些事件与确认事件流连接起来
  • 将连接的 DataFrame 写入 HBase 表。

我在 spark-shell 中测试了代码,它工作正常,低于伪代码(我使用的是 Scala)。

val tableA = spark.table("tableA")

val startingOffset = "earliest"

val streamOfData = .readStream 
  .format("kafka") 
  .option("startingOffsets", startingOffsets)
  .option("otherOptions", otherOptions)

val joinTableAWithStreamOfData = streamOfData.join(tableA, Seq("a"), "inner")

joinTableAWithStreamOfData 
  .writeStream
  .foreach(
    writeDataToHBaseTable()
  ).start()
  .awaitTermination()

现在我想安排此代码定期运行,例如每 15 分钟一次,我正在努力理解如何在此处使用检查点

在每次运行此代码时,我想仅从流中读取我在上一次运行中尚未读取的事件,并将这些新事件与我的日志表内部连接,所以只将新数据写入最终的 HBase 表。

我在 HDFS 中创建了一个目录来存储检查点文件。 我将该位置提供给我用来调用 spark 代码的 spark-submit 命令。

spark-submit --conf spark.sql.streaming.checkpointLocation=path_to_hdfs_checkpoint_directory 
--all_the_other_settings_and_libraries

此时代码每 15 分钟运行一次,没有任何错误,但它基本上没有做任何事情,因为它没有将新事件转储到 HBase 表。 检查点目录也是空的,而我假设必须在那里写入一些文件?

是否需要调整 readStream 函数才能从最新的检查点开始读取?

val streamOfData = .readStream 
  .format("kafka") 
  .option("startingOffsets", startingOffsets) ??
  .option("otherOptions", otherOptions)

我真的很难理解有关此的 spark 文档。

提前谢谢你!

【问题讨论】:

  • 你能发布在 scala shell 中工作的完整代码或实际代码吗??
  • 我的猜测是 dataframe joinTableAWithStreamOfData 是空的,因为这个 writestream 没有被触发或启动。如果它被启动,就会创建检查点位置。
  • 感谢您的 cmets! @Rayan 我已经看到了那个答案,这促使我写下我的问题:) 所以应该在某处 writing 流时设置检查点?我想知道这实际上是如何工作的,因为我不仅将流写入表,而且首先我将流与其他数据连接起来,然后写入 HBase,所以我想知道 writeStream 上的检查点是如何工作的

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


【解决方案1】:

触发器

“现在我想安排此代码定期运行,例如每 15 分钟一次,我正在努力理解如何在这里使用检查点。

如果您希望每 15 分钟触发一次工作,您可以使用Triggers

您不需要专门“使用”检查点,只需提供可靠的(例如 HDFS)检查点位置,见下文。

检查点

在每次运行此代码时,我只想从流中读取我在上一次运行中尚未读取的事件 [...]"

在 Spark 结构化流应用程序中从 Kafka 读取数据时,最好将检查点位置直接设置在您的 StreamingQuery 中。 Spark 使用此位置创建检查点文件,以跟踪应用程序的状态并记录已从 Kafka 读取的偏移量。

重新启动应用程序时,它将检查这些检查点文件以了解从何处继续从 Kafka 读取,因此它不会跳过或错过任何消息。您无需手动设置startingOffset。

请务必记住,仅允许对应用程序代码进行特定更改,以便检查点文件可用于安全重启。可以在Recovery Semantics after Changes in a Streaming Query 上的结构化流编程指南中找到一个很好的概述。


总体而言,对于从 Kafka 读取的高效 Spark Structured Streaming 应用程序,我推荐以下结构:

val spark = SparkSession.builder().[...].getOrCreate()

val streamOfData = spark.readStream 
  .format("kafka") 
// option startingOffsets is only relevant for the very first time this application is running. After that, checkpoint files are being used.
  .option("startingOffsets", startingOffsets) 
  .option("otherOptions", otherOptions)
  .load()

// perform any kind of transformations on streaming DataFrames
val processedStreamOfData = streamOfData.[...]


val streamingQuery = processedStreamOfData 
  .writeStream
  .foreach(
    writeDataToHBaseTable()
  )
  .option("checkpointLocation", "/path/to/checkpoint/dir/in/hdfs/"
  .trigger(Trigger.ProcessingTime("15 minutes"))
  .start()

streamingQuery.awaitTermination()

【讨论】:

  • 我正在尝试在 pyspark 中进行相同的实现。我担心的是我想确保应用程序是否重新启动,然后必须从 Kafka 分区的相同偏移量启动。我在 readStream() 中指定了 checkPointingLocation 选项,但它似乎不起作用(根据日志,它指定了 tmp 文件夹中的另一个位置)。所以我想知道在这个模型中为什么我们在 writeStream() 中指定检查点位置,不应该在 readStream() 中吗?
  • 据我了解,checkpoint 会记录处理了哪些偏移量(已经写入成功),所以放在 writeStream() 中应该更合适。
猜你喜欢
  • 2021-06-03
  • 2020-09-08
  • 2022-01-14
  • 2020-09-12
  • 2021-01-17
  • 2018-07-12
  • 2020-03-19
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多