【发布时间】:2020-10-09 12:38:12
【问题描述】:
我们在 Kubernetes 上使用 Apache Flink job cluster,它由一个 Job Manager 和两个 Task Managers 组成,每个 Task Managers 有两个插槽。使用Lightbend Cloudflow框架部署和配置集群。
我们还使用RocksDB 状态后端和兼容 S3 的存储来实现持久性。考虑从 CLI 创建 savepoints 没有任何问题。我们的工作由几个键控状态 (MapState) 组成,而且往往相当庞大(我们预计每个状态至少 150 Gb)。作业的Restart Strategy 设置为Failure Rate。我们在工作中使用Apache Kafka 作为源和汇。
我们目前正在进行一些测试(主要是 PoC),但仍有一些问题悬而未决:
我们进行了一些综合测试,并将不正确的事件传递给工作。这导致Exceptions 在执行期间被抛出。由于Failure Rate 策略,会发生以下步骤:通过源读取来自 Kafka 的损坏消息 -> 操作员尝试处理该事件并最终抛出 Exception -> 作业重新启动并读取 THE SAME em> 来自 Kafka 的记录,如上一步 -> 操作员失败 -> Failure Rate 最终超过给定值,作业最终停止。接下来我该怎么办?如果我们尝试重新启动作业,似乎它将使用最新的 Kafka 消费者状态恢复并再次读取损坏的消息,从而导致我们回到前面提到的行为?哪些是解决此类问题的正确步骤? Flink 是否使用了任何一种所谓的Dead Letter Queues?
另一个问题是关于检查点和恢复机制。我们目前无法弄清楚在作业执行期间引发的哪些异常被认为是关键的,并导致作业失败,然后从最新的检查点自动恢复?正如在前一个案例中所描述的,在作业中引发的普通Exception 会导致持续重新启动,最终导致作业终止。当我们的集群发生某些事情(Job Manager 失败,Task Manager 失败或其他事情)时,我们正在寻找一个案例来重现,从而导致从最新的检查点自动恢复。考虑到 Kubernetes 集群中的这种情况,欢迎提出任何建议。
我们查阅了 Flink 官方文档,但没有找到任何相关信息或可能以错误的方式感知它。非常感谢!
【问题讨论】:
标签: kubernetes apache-kafka apache-flink