【问题标题】:Two questions on Flink externalized checkpointsFlink外部化检查点的两个问题
【发布时间】:2018-06-07 09:35:19
【问题描述】:

我有两个关于 Flink 外部化检查点的问题

(Q1) 我可以在 flink-conf.yaml 中设置“state.checkpoints.dir”来让外部化的检查点正常工作,但是当我从 IDE 运行 flink 时如何实现相同的效果呢?我尝试了 (http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/state-checkpoints-dir-td17921.html) 中提到的 GlobalConfiguration 方法,但没有运气。我就是这样做的:

Configuration cfg =
                GlobalConfiguration.loadConfiguration();
cfg.setString("state.checkpoints.dir", "file:///tmp/checkpoints/state");
env.getCheckpointConfig().enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

这是 IDE 中显示的错误消息:

Exception in thread "main" org.apache.flink.runtime.client.JobExecutionException: Failed to submit job ef7050e2308a4787d983d80f3c07f55c (Long Taxi Rides (checkpointed))
    at org.apache.flink.runtime.jobmanager.JobManager.org$apache$flink$runtime$jobmanager$JobManager$$submitJob(JobManager.scala:1325)
    at org.apache.flink.runtime.jobmanager.JobManager$$anonfun$handleMessage$1.applyOrElse(JobManager.scala:447)
    at scala.runtime.AbstractPartialFunction.apply(AbstractPartialFunction.scala:36)
    at org.apache.flink.runtime.LeaderSessionMessageFilter$$anonfun$receive$1.applyOrElse(LeaderSessionMessageFilter.scala:38)
    at scala.runtime.AbstractPartialFunction.apply(AbstractPartialFunction.scala:36)
    at org.apache.flink.runtime.LogMessages$$anon$1.apply(LogMessages.scala:33)
    at org.apache.flink.runtime.LogMessages$$anon$1.apply(LogMessages.scala:28)
    at scala.PartialFunction$class.applyOrElse(PartialFunction.scala:123)
    at org.apache.flink.runtime.LogMessages$$anon$1.applyOrElse(LogMessages.scala:28)
    at akka.actor.Actor$class.aroundReceive(Actor.scala:502)
    at org.apache.flink.runtime.jobmanager.JobManager.aroundReceive(JobManager.scala:122)
    at akka.actor.ActorCell.receiveMessage(ActorCell.scala:526)
    at akka.actor.ActorCell.invoke(ActorCell.scala:495)
    at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:257)
    at akka.dispatch.Mailbox.run(Mailbox.scala:224)
    at akka.dispatch.Mailbox.exec(Mailbox.scala:234)
    at scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
    at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
    at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
    at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)
Caused by: java.lang.IllegalStateException: CheckpointConfig says to persist periodic checkpoints, but no checkpoint directory has been configured. You can configure configure one via key 'state.checkpoints.dir'.
    at org.apache.flink.runtime.checkpoint.CheckpointCoordinator.<init>(CheckpointCoordinator.java:211)
    at org.apache.flink.runtime.executiongraph.ExecutionGraph.enableCheckpointing(ExecutionGraph.java:478)
    at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilder.buildGraph(ExecutionGraphBuilder.java:291)
    at org.apache.flink.runtime.jobmanager.JobManager.org$apache$flink$runtime$jobmanager$JobManager$$submitJob(JobManager.scala:1277)
    ... 19 more

Process finished with exit code 1

(Q2)在检查点的文档(https://ci.apache.org/projects/flink/flink-docs-release-1.4/dev/stream/state/checkpointing.html)中,它说“这样,如果您的工作失败,您将有一个检查点可以从周围恢复。”,取消的工作怎么样?新作业会继续现有的检查点还是从自己的检查点开始?

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    您可以控制在取消作业时是否删除外部化检查点。如果你想保留它们,你可以这样做:

    CheckpointConfig config = env.getCheckpointConfig();
    config.enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
    

    有关详细信息,请参阅the docs

    resume from an externalized checkpoint 这样做(与从保存点恢复相同):

    $ bin/flink run -s :checkpointMetaDataPath [:runArgs]
    

    【讨论】:

    • 谢谢 Dave - 我忘记了在取消时删除检查点只是默认行为,可以进行不同的配置。
    • 当我从 flink ui 取消并再次启动我的应用程序时,flink 是否会自动添加此 checkpointMetaDataPath ?从命令恢复时,这应该可以工作: bin/flink run -s :checkpointMetaDataPath [:runArgs]
    • @David Anderson 我们有办法从存储在 S3 中的检查点自动恢复吗?这变得很棘手,因为 -s :checkpointMetaDataPath 需要有 flink 的 job_id 并在 tun 时找出这个 job id?
    • @UmeshK Flink 通常会自动从检查点恢复。从保存点或外部检查点开始不需要 job_id。所以我不明白你的问题——你能问一个新问题,并提供更多关于你想要做什么的细节吗?
    【解决方案2】:

    第一个问题,使用自定义配置创建本地环境:

    val conf = new Configuration()
    conf.setString(CoreOptions.CHECKPOINTS_DIRECTORY, "file:///user/flink/checkpoint/storing/")
    StreamExecutionEnvironment.createLocalEnvironment(4, conf)
    

    第二个问题,正如大卫所说:

    config.enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
    

    【讨论】:

      【解决方案3】:

      从 Eclipse 重新设置检查点目录,通常我只是在设置要使用的后端时这样做,例如

      env.setStateBackend(new FsStateBackend(options.getCheckpointDir()));
      

      重新取消作业 - 检查点目录被删除。如果您想在停止(取消)您的工作后从已知状态恢复,您需要创建一个保存点。

      【讨论】:

      • options 只是我自己的带有命令行参数值的类。您可以将任何有效路径(从 Eclipse 运行时通常为 file:///xxx)传递给 FsStateBackend 构造函数。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-04-04
      • 2023-01-18
      • 2021-02-15
      • 2020-09-16
      相关资源
      最近更新 更多