【发布时间】:2020-02-27 07:39:14
【问题描述】:
我们的团队在我们的 K8S 集群中搭建了一个 Flink Session Cluster。我们选择了 Flink Session Cluster 而不是 Job Cluster,因为我们有许多不同的 Flink Jobs,所以我们想将 Flink 的开发和部署与我们的 Job 的开发和部署解耦。我们的 Flink 设置包含:
- 单个 JobManager 作为 K8S pod,无高可用性 (HA) 设置
- 多个TaskManager,每个都作为一个K8S pod
我们在单独的存储库中开发我们的作业,并在合并代码时部署到 Flink 集群。
现在,我们注意到 JobManager 作为 K8S 中的 pod 可以随时被 K8S 重新部署。因此,一旦重新部署,它就会失去所有工作。为了解决这个问题,我们开发了一个脚本来持续监控 Flink 中的作业,如果作业没有运行,脚本会重新提交作业到集群。由于脚本可能需要一些时间来发现并重新提交作业,因此经常会出现小的服务中断,我们正在考虑是否可以改进。
到目前为止,我们有一些想法或问题:
一种可能的解决方案是:当 JobManager 被(重新)部署时,它将获取最新的 Jobs jar 并运行作业。这个解决方案看起来总体不错。尽管如此,由于我们的作业是在单独的 repo 中开发的,因此我们需要一个解决方案让集群在作业发生更改时注意到最新的作业,要么 JobManager 不断轮询最新的作业 jar,要么 Jobs repo 部署最新的作业 jar。
我看到 Flink HA 功能可以存储检查点/保存点,但不确定 Flink HA 是否已经可以处理这个重新部署问题?
有人对此有任何意见或建议吗?谢谢!
【问题讨论】: