【发布时间】:2020-09-16 09:07:50
【问题描述】:
我们正在运行一个 ListState 介于 300GB 和 400GB 之间的作业,有时该列表可能会增长到数千。在我们的用例中,每个项目都必须有自己的 TTL,因此我们为这个 ListState 的每个新项目创建一个新的 Timer,并在 S3 上使用 RocksDB 后端。
目前大约有 140+ 百万个计时器(将在 event.timestamp + 40 天 触发)。
我们的问题是,作业的检查点突然卡住了,或者非常慢(比如几个小时内 1%),直到最终超时。它通常会在一段非常简单的代码上停止(flink 仪表板显示0/12 (0%),而前几行显示12/12 (100%)):
[...]
val myStream = env.addSource(someKafkaConsumer)
.rebalance
.map(new CounterMapFunction[ControlGroup]("source.kafkaconsumer"))
.uid("src_kafka_stream")
.name("some_name")
myStream.process(new MonitoringProcessFunction()).uid("monitoring_uuid").name(monitoring_name)
.getSideOutput(outputTag)
.keyBy(_.name)
.addSink(sink)
[...]
更多信息:
- AT_LEAST_ONCE 检查点模式似乎比 EXACTLY_ONCE 更容易卡住
- 几个月前,该州的数据量达到了 1.5TB,我认为数十亿个计时器没有任何问题。
- 运行两个任务管理器的机器上的 RAM、CPU 和网络看起来正常
state.backend.rocksdb.thread.num = 4- 第一个事件发生在我们收到大量事件(大约数百万分钟)时,但不是在前一个事件中发生。
- 所有事件都来自 Kafka 主题。
- 在 AT_LEAST_ONCE 检查点模式下,作业仍然正常运行和消耗。
这是我们第二次遇到拓扑运行良好,每天有几百万个事件并突然停止检查点。我们不知道是什么原因造成的。
任何人都可以想到什么会突然导致检查点卡住?
【问题讨论】:
标签: apache-flink flink-streaming rocksdb