【发布时间】:2021-06-04 18:39:22
【问题描述】:
我们有一个 Flink 作业(Flink 版本:1.9),它通过 key 连接两个 kafka 源,对于每个 key,启动一个 5 分钟的计时器,消息在 Flink 状态下缓存,当计时器结束时,合并相同的消息将密钥(通常每个密钥有 1~5 条消息)转换为胖消息并将其发送到 kafka。
两个kafka来源:
- source1(160 个分区,每分钟 20~3000 万条消息),
- source2(30 个分区,每分钟 1~3 百万条消息)。
平面图只是反序列化了 kafka 消息。
KeyedProcess 是计时器和 Flink 状态发挥作用的地方。
我已经尝试了一些来提高性能,例如取模关键是减少定时器的数量,或者增加硬件(目前2000c 4000gb),或者调整算子的并行度。
目前的问题是,当 source1 超过每分钟 2500 万条消息时,消耗速度会急剧下降,并且永远不会恢复。如果低于 2500 万条消息/分钟,它可以正常工作。
kafka 集群本身似乎没有问题,因为有另一个系统正在读取它,并且该系统没有任何消耗速度问题。
有人能解释一下吗?如何解决原因?或者我可以尝试什么?添加更多硬件是个好主意吗(我认为 2000c&4000gb 已经是一个巨大的资源量)?非常感谢。
【问题讨论】:
-
状态后端是如何配置的?
-
@DavidAnderson 嗨,大卫,我使用 RocksDB 状态后端。顺便说一句,检查点已禁用。
标签: java apache-kafka apache-flink