【问题标题】:Flink checkpointing working for ProcessFunction but not for AsyncFunctionFlink 检查点适用于 ProcessFunction 但不适用于 AsyncFunction
【发布时间】:2021-12-26 20:44:49
【问题描述】:

我为ProcessFunction 操作员启用了操作员检查点并顺利工作。

在作业失败时,我可以看到操作员状态如何在 snapshotState() 挂钩上外部化,在恢复时,我可以看到状态如何在 initializeState() 挂钩上恢复。

但是,当我尝试在 AsyncFunction 上实现 CheckpointedFunction 接口和上述两种方法时,它似乎不起作用。我所做的几乎与ProcessFunction 相同......但是当工作在失败后关闭时,它似乎并没有被snapshotState() 钩子停止,并且在工作恢复时,context.isRestored() 总是假的。

为什么CheckpointedFunction.snapshotState()CheckpointedFunction.initializeState() 不与AsyncFunction 一起执行,而与ProcessFunction 一起执行?

已编辑: 出于某种原因,我的检查站需要很长时间。我相信我的配置非常标准,1 秒的间隔,500 毫秒的最小暂停,恰好一次。没有其他调整。
我从检查点协调器那里得到了这些痕迹

o.a.f.s.r.t.SubtaskCheckpointCoordinatorImpl - Time from receiving all checkpoint barriers/RPC to executing it exceeded threshold: 93905ms
2021-11-23 16:25:01 INFO  o.a.f.r.c.CheckpointCoordinator - Completed checkpoint 4 for job 239d7967eac7900b33d7eadd483c9447 (671604 bytes in 112071 ms).

如果我尝试设置 checkpointTimeout,我需要按顺序或 5 分钟左右设置一些内容。这么小的状态(就是一个Counter和一个Long)的checkpoint怎么要5分钟?

我还读到 NFS 卷是一个麻烦的秘诀,但到目前为止我还没有在集群上运行它,我只是在我的本地文件系统上测试它

【问题讨论】:

    标签: apache-flink checkpointing


    【解决方案1】:

    AsyncFunction 根本不支持状态。原因是状态原语不同步,因此会在AsyncFunction 中产生不正确的结果。这与没有KeyedAsyncFunction 的原因相同。

    如果 Flink 实现了 https://cwiki.apache.org/confluence/display/FLINK/FLIP-22%3A+Eager+State+Declaration,那么它可以简单地在每个异步调用上附加状态并在异步成功时更新。

    您可以在限制周围使用链式地图和插槽共享组做一些诡计,但它相当hacky。

    【讨论】:

    • this stackoverflow.com/a/62476382/11217621 似乎暗示尽管 AsyncFunction 不允许保持任何键控状态,但仍然可以与 CheckpointedFunction 混合......这是不正确的吗?我需要在作业失败之前保留最后一个元素,以便我的异步功能可以从那里恢复。有没有办法做到这一点?
    • 我仔细检查了当前的代码库,你是对的。应该支持CheckpointedFunction。您可以在 IDE 中执行您的代码(始终推荐)并检查 this method 是否在检查点上执行了哪些操作?
    • 您能否详细说明我在哪里以及如何连接 StreamingFunctionUtils.snapshotFunctionState() 以及我从哪里得到它的参数? ...在 AsyncFuntion 的 initializeState() 上?每次通过 asyncInvoke()?你可以用一些代码来说明它吗?
    • 对不起,我的意思是在你的 IDE 中执行 Flink 时在你的 IDE 中的这个函数处设置一个断点(或者你将一个远程调试器附加到你的任务管理器)。你没有办法使用它。它在 Flink 内部使用。
    • 我认为这是一个不相关的问题,您可以在 Flink 邮件列表中询问。我的第一个猜测是您的管道中有背压(可能来自 asyncIO),因此检查点屏障需要 100 秒才能通过。
    猜你喜欢
    • 2022-11-03
    • 2021-07-12
    • 1970-01-01
    • 1970-01-01
    • 2011-05-20
    • 2020-04-13
    • 2021-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多