【发布时间】:2020-03-10 09:20:34
【问题描述】:
目前,我们正在尝试弄清楚如何有效地使用 Flink,并且我们仍在尝试理解一切。
我们在一个独立集群上运行了大约 60 个真正轻量级的作业,这在一个普通的 EC2 实例上运行良好。但是,一旦我使用本地 RocksDB 状态后端启用检查点,集群就会以意想不到的方式运行,停止作业,尝试重新启动它们只是丢弃所有它们并将错误日志留空。之后,在 Flink 中没有留下任何作业或 jar 的痕迹。
我知道,对于每个作业,都会保留 JobManager 总内存的一小部分,同样,对于每个作业,本地 RocksDB 都会在同一台机器上实例化,但我认为它们同样轻巧,不需要太多内存/CPU 容量.与之前的稳定集群相比,只需添加行 env.enableCheckpointing(1000); 就会导致一切完全失败。
我个人认为我们可能已经达到了我们独立 Flink 集群的极限,即使增加内存也不够了,但我想在开始构建分布式 Flink 集群之前确认这一点(我们需要自动化一切,这就是我现在犹豫的原因)。我不确定是否例如将 RocksDB 检查点存储在像 S3 这样的专用存储单元中甚至可以解决这个问题,如果资源消耗(硬盘除外)会受到影响。
迁移到分布式环境是解决我们问题的唯一方法,还是这表明存在其他问题,可以通过适当的配置来解决?
编辑:也许我应该补充一点,还没有加载,我们还没有谈论传入的数据,但关于作业仍在运行。目前 FlinkSources 中只有 100 条记录,但我们甚至无法达到正在处理的程度。
编辑2:
这部分一直是作业代码的一部分:
try {
env.setStateBackend((StateBackend) new RocksDBStateBackend("file://" + "/somePathOnInstance"));
} catch (IOException e1) {
// TODO Auto-generated catch block
e1.printStackTrace();
}
我们添加了以下代码行:
env.enableCheckpointing(CHECKPOINTING_INTERVAL_MS);
不需要对 StateBackend 进行类型转换,因为根据文档,RocksDBStateBackend 类的 1.9.1 版本应该已经实现了 StateBackend 而不是 AbstractStateBackend。但是,该文档与我们从 Maven 获得的实际类不同,所以就是这样。
【问题讨论】:
-
您是否同时启用了检查点并切换到使用 RocksDB,或者您之前是否成功使用过 RocksDB,但没有检查点?
-
@DavidAnderson 我们之前已经启用了 RocksDB,我在最初的帖子中添加了详细信息,请查看最近的编辑。但是,如果不使用检查点,RocksDB 可能只是用于捕获状态,在我们的例子中它非常小(部分只是一个布尔值)。我只能假设它比使用实际检查点占用更少的资源?
-
开启检查点确实会增加资源需求。您真的在单个 EC2 实例上同时运行 60 个作业吗?你同时为所有这些都打开了检查点?
-
@DavidAnderson 从这个问题来看,我认为这很不寻常?答案是肯定的,60 个工作(主要是 1 个源、1 个接收器和一个介于两者之间的简单函数),这实际上在 t3 上没什么大问题。中等实例没有检查点。然而,没有繁重的负载,最多只有几千条记录,但它是这样工作的。
-
我不知道这是否会有所改进,但我只是指出可以在同一个作业中运行多个管道。使用这种技术来减少工作的总数应该会减少总体资源需求。
标签: apache-flink