我们是否需要在 Spark Structured Streaming 中同时检查 Kafka 的 readStream 和 writeStream?

use*_*400 2 apache-spark spark-streaming

我们是否需要在 Spark Structured Streaming 中检查 Kafka 的 readStream 和 writeStream ?我们什么时候需要检查这两个流或其中一个流?

Yur*_*ruk 5

需要检查点来保存有关流处理数据的信息,并且在失败的情况下,火花可以从上次保存的进度点恢复。处理意味着它从源读取,(转换)并最终写入接收器。

因此,无需分别为 reader 和 writer 设置检查点,因为在恢复后不处理仅读取但未写入 sink 的数据是没有意义的。此外,检查点位置只能设置为 DataStreamWriter 的一个选项(从 返回dataset.writeStream())并且在启动流之前。

这是带有检查点的简单结构化流的示例:

session
    .readStream()
    .schema(RecordSchema.fromClass(TestRecord.class))
    .csv("s3://test-bucket/input")
    .as(Encoders.bean(TestRecord.class))
    .writeStream()
    .outputMode(OutputMode.Append())
    .format("csv")
    .option("path", "s3://test-bucket/output")
    .option("checkpointLocation", "s3://test-bucket/checkpoint")
    .queryName("test-query")
    .start();
Run Code Online (Sandbox Code Playgroud)

  • 这是否意味着,当我们初始化“readStream”时,不需要设置 Kafka 相关选项,如“startingOffsets”、“auto.offset.reset”、“enable.auto.commit”等?并且“readStream”将仅检索从“checkpointLocation”中保存的偏移量开始的消息? (8认同)