【问题标题】:Flink reads Kafka, consume speed drops dramatically in some caseFlink 读取 Kafka,某些情况下消耗速度急剧下降
【发布时间】:2021-06-04 18:39:22
【问题描述】:

我们有一个 Flink 作业(Flink 版本:1.9),它通过 key 连接两个 kafka 源,对于每个 key,启动一个 5 分钟的计时器,消息在 Flink 状态下缓存,当计时器结束时,合并相同的消息将密钥(通常每个密钥有 1~5 条消息)转换为胖消息并将其发送到 kafka。

两个kafka来源:

  1. source1(160 个分区,每分钟 20~3000 万条消息),
  2. source2(30 个分区,每分钟 1~3 百万条消息)。

平面图只是反序列化了 kafka 消息。

KeyedProcess 是计时器和 Flink 状态发挥作用的地方。

我已经尝试了一些来提高性能,例如取模关键是减少定时器的数量,或者增加硬件(目前2000c 4000gb),或者调整算子的并行度。

目前的问题是,当 source1 超过每分钟 2500 万条消息时,消耗速度会急剧下降,并且永远不会恢复。如果低于 2500 万条消息/分钟,它可以正常工作。

kafka 集群本身似乎没有问题,因为有另一个系统正在读取它,并且该系统没有任何消耗速度问题。

有人能解释一下吗?如何解决原因?或者我可以尝试什么?添加更多硬件是个好主意吗(我认为 2000c&4000gb 已经是一个巨大的资源量)?非常感谢。

【问题讨论】:

  • 状态后端是如何配置的?
  • @DavidAnderson 嗨,大卫,我使用 RocksDB 状态后端。顺便说一句,检查点已禁用。

标签: java apache-kafka apache-flink


【解决方案1】:

您可以先附加一个分析器来查看瓶颈在哪里。 (也许是磁盘?)

RocksDB 在某种程度上表现不佳似乎是合理的。 可能需要进行一些调整。您应该能够通过enabling the RocksDB native metrics 获得一些见解,并查看问题发生时各种 RocksDB 指标如何变化。

这些是一些更有用的指标:

estimate-live-data-size
estimate-num-keys
num-running-compactions
num-live-versions
estimate-pending-compaction-bytes
num-running-flushes
size-all-mem-tables
block-cache-usage

根据此工作负载的运行位置和方式,您可能会遇到某种速率限制或节流。一个有趣的例子见The Impact of Disks on RocksDB State Backend in Flink: A Case Study

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-01-29
    • 1970-01-01
    • 2015-12-24
    • 1970-01-01
    • 2012-04-01
    • 1970-01-01
    相关资源
    最近更新 更多