【问题标题】:How to restore state after a restart from chosen source (not necessarily last checkpoint)从所选源重新启动后如何恢复状态(不一定是最后一个检查点)
【发布时间】:2019-09-11 03:26:05
【问题描述】:

我一直在尝试从之前的检查点重新启动我的 Apache Flink,但运气不佳。我已经将代码上传到 GitHub,这里是主要类: https://github.com/edu05/wordcount/blob/restart/src/main/java/edu/streaming/AppWithKafka.java

这是一个简单的字数统计程序,只是我希望程序在重新启动后继续使用它已经计算的字数。

我已经阅读了文档并尝试了一些东西,但一定是缺少一些愚蠢的东西,有人可以帮忙吗?

另外:最终目标是将 wordcount 程序的输出生成到压缩的 kafka 主题中,我将如何通过首先使用压缩的主题来加载应用程序的状态,在这种情况下,它既可以作为输出以及程序的检查点机制?

非常感谢

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    Flink 的检查点用于在失败后自动重启。如果要手动重新启动,请使用savepointexternalized checkpoint

    如果您已经尝试过,但仍然遇到问题,请提供有关您尝试过的更多详细信息。

    【讨论】:

    • 嗨,大卫,感谢您的关注。我已经更新了分支,尝试以预加载状态启动我的应用程序,其中在处理任何新消息之前单词“mario”的计数为 3(我实现了 CheckpointedFunction 接口)但它会引发异常。原始用例更复杂,但让这个更简单的场景工作就足够了。我们可能也有不同的重启概念。对我来说,重启会导致任意数量的节点(可能全部)宕机,偶尔会发生,除了重启机器之外不需要人工干预。
    • 顺便说一句,即使之前提交的代码也不会从它停止的地方重新开始,所以每次我启动应用程序时都会重新计算字数:(
    • 我们可以重新回答这个问题吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-09-07
    • 2020-06-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多