【问题标题】:How to share state between runs of streaming jobs?如何在流作业运行之间共享状态?
【发布时间】:2023-01-18 21:09:21
【问题描述】:

由于业务需求,我每天使用 Trigger.Once 方法触发一个 Spark 流作业。

StreamingQuery query = joinedDf
                       .writeStream()
                       .outputMode("append")
                       .format("parquet")
                       .option("path", resultPath)
                       .option("checkpointLocation", checkpointLocationPathForDate)
                       .trigger(Trigger.Once())
                       .start();

我正在使用 map flatMapGroupsWithState 以便我们可以存储分组数据的状态 (GroupState)。 我在某个地方读到每个 StreamingQuery 的 checkpointLocation 应该不同。因此我使用这样的检查点位置:/path/to/nfs/checkpoint/<current date in format: yyyyMMdd>

每天,Spark 作业都会处理文件夹/path/to/data/<current date in format: yyyyMMdd> 中的文件

我想访问昨天 Spark 作业的状态,因为昨天的数据可能包含今天数据所需的相关状态。

但是,Spark 将状态数据存储在 checkpointLocation 即 /path/to/nfs/checkpoint/<current date in format: yyyyMMdd>/<queryName>/state 中,因此当使用不同的 checkpointLocation 时,无法访问它。

那么,如何访问存储在先前 Spark 作业的 checkpointLocation 中的 GroupState 数据?可以对不同的 StreamingQueries 使用相同的 checkpointLocation 吗?

编辑: 我尝试对昨天的 StreamingQuery 和今天的 StreamingQuery 使用相同的 checkpointLocation,并且 Spark 恢复了昨天批次的状态,这是我想要的,但是这在任何地方都有记录吗?当在每日批次之间使用相同的 checkpointLocation 时,这是预期的行为还是可能的行为不端?

【问题讨论】:

    标签: apache-spark spark-structured-streaming


    【解决方案1】:

    如何访问存储在先前 Spark 作业的 checkpointLocation 中的 GroupState 数据?

    你不应该。从技术上讲,您可以(通过一些额外的编码)但是您应该考虑其他查询特有的许多事情(例如,有状态的操作员 ID)。使用风险自负。

    可以对不同的 StreamingQueries 使用相同的 checkpointLocation 吗?

    不,您不应该在不同的流式查询之间共享相同的checkpointLocation。一是他们与他们的运营商不同,所以数字可能不匹配,即使他们匹配,接收器也可能不同,因此一些数据可能会被跳过(已经处理过)。

    我尝试对昨天的 StreamingQuery 和今天的 StreamingQuery 使用相同的 checkpointLocation,并且 Spark 恢复了昨天批次的状态,这是我想要的,但是这在任何地方都有记录吗?当在每日批次之间使用相同的 checkpointLocation 时,这是预期的行为还是可能的行为不端?

    这已记录在案,这正是 checkpointLocation 应该如何工作。它是具有给定时间的流式查询状态的目录。

    引用Recovering from Failures with Checkpointing

    在发生故障或故意关闭的情况下,您可以恢复先前查询的先前进度和状态,并从中断处继续。这是使用检查点和预写日志完成的。您可以使用检查点位置配置查询,查询会将所有进度信息(即每个触发器中处理的偏移量范围)和运行聚合(例如快速示例中的字数)保存到检查点位置。此检查点位置必须是 HDFS 兼容文件系统中的路径,并且可以在启动查询时设置为 DataStreamWriter 中的选项。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-07-27
      • 2016-11-02
      • 2021-03-27
      • 2010-11-12
      • 1970-01-01
      • 2019-05-07
      • 1970-01-01
      相关资源
      最近更新 更多