【问题标题】:Spark Streaming - Restarting from checkpoint replays last batchSpark Streaming - 从检查点重新开始重播最后一批
【发布时间】:2017-05-31 14:56:04
【问题描述】:

我们正在尝试构建一个容错的火花流作业,我们遇到了一个问题。这是我们的场景:

    1) Start a spark streaming process that runs batches of 2 mins
    2) We have checkpoint enabled. Also the streaming context is configured to either create a new context or build from checkpoint if one exists
    3) After a particular batch completes, the spark streaming job is manually killed using yarn application -kill (basically mimicking a sudden failure) 
    4) The spark streaming job is then restarted from checkpoint

我们遇到的问题是,在重新启动 spark 流作业后,它会重播最后一个成功的批处理。它总是这样做,只是重播最后一个成功的批次,而不是之前的批次

这样做的副作用是该批次的数据部分是重复的。我们甚至尝试在最后一个成功的批处理之后等待超过一分钟,然后再终止进程(以防写入检查点需要时间),但这没有帮助

有什么见解吗?我没有在这里添加代码,希望有人也遇到过这个问题并可以提供一些想法或见解。如果有帮助,也可以发布相关代码。不应该在批处理成功后立即触发流检查点,以便在重新启动后不会重播?我将 ssc.checkpoint 命令放在哪里重要吗?

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    您在问题的最后一行有答案。 ssc.checkpoint() 的位置很重要。当您使用保存的检查点重新启动作业时,该作业会出现正在保存的任何内容。因此,在您的情况下,当您在批处理完成后终止作业时,最近的作业是最后一个成功的作业。到了这个时候,你可能已经明白检查点主要是从你离开的地方重新开始——尤其是对于失败的工作。

    【讨论】:

    • 感谢拉姆齐。 “到这个时候,你可能已经明白检查点主要是从你离开的地方重新开始——尤其是对于失败的工作”。是的,这是有道理的。这是我没有得到的,当我杀死火花工作时,最后一个状态是成功的批处理。为什么在重新启动时会重播该批次?这是一个限制还是需要改变我设置检查点的方式?
    • 理想情况下,流媒体应该连续运行。因此,如果在工作之间发生了一些事情,检查点应该确保它在它停止的地方被拾取。关于重复 - 这可以由 checkpoint() 方法的位置以及如何处理下游应用程序中的重复来决定。最后,您将因失败而批量丢失一些数据或重做整个批次。但这些情况很少见,因为我们谈论的是流媒体。您的情况很棘手,因为您使用流式传输来频繁停止和启动
    【解决方案2】:

    有两件事需要注意。

    1] 确保在重新启动程序时在 getOrCreate 流式上下文方法中使用相同的检查点目录。

    2] 将“spark.streaming.stopGracefullyOnShutdown”设置为“true”。这允许spark完成对当前数据的处理并相应地更新检查点目录。如果设置为false,可能会导致检查点目录中的数据损坏。

    注意:如果可能,请发布代码 sn-ps。是的,ssc.checkpoint 的位置确实很重要。

    【讨论】:

      【解决方案3】:

      在这种情况下,应确保 Spark 应用程序重启后流式上下文方法中使用的检查点目录相同。希望它会有所帮助

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-10-09
        相关资源
        最近更新 更多