【发布时间】: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