【问题标题】:How to make lag go down in kafka stream based application如何在基于 kafka 流的应用程序中降低延迟
【发布时间】:2019-10-08 11:27:12
【问题描述】:

我有一个带有 3 个 kafka 机器集群的真实环境,它正在接收大量数据。每个主题有 25 个分区,复制因子设置为 2。

我的应用程序(基于 kafka 流的应用程序)从这个 kafka 集群获取数据已停机一个多月。现在,每个分区都有大量的滞后;达到 90000000。

我知道以下参数:

max.poll.records ; default —> 500
max.partition.fetch.bytes ; default —> 1048576
fetch.max.bytes ; default —> 52428800
fetch.min.bytes ; default —> 1

max.poll.interval.ms ; default —> 300000
request.timeout.ms; default —> 30000
session.timeout.ms ; default —> 10000

我有 2 个消费者节点(使用来自 kafka 集群的数据的相同组 ID)。

但是,它并没有赶上滞后,它保持不变。谁能建议如何改进以降低延迟?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    如果您的应用程序宕机了一个月,一些记录已过期,因为主题中的默认保留期为 7 天,因此您很可能丢失了一些消息。此外,默认偏移重置保留 1 天或 7 天,具体取决于您的 Kafka Streams 版本。似乎您有auto.offset.reset: earliest,因此它从每个分区的开头开始使用消息。如果您需要跳过所有消息并仅使用新消息,则应设置 auto.offset.reset: latest 并将 application.id 值更改为新值。

    如果您想并行消费消息并加快延迟减少,您可以将配置num.stream.threads 设置为12 之类的值(num.stream.threads * numberOfConsumerNodes 应该小于或等于numberOfPartitions,否则,一些线程将处于空闲状态),或者需要增加消费者节点的数量。

    【讨论】:

    • 我每个主题有 25 个分区,我有 25 个消费者,所以我已经照顾好了。还有其他因素可以帮助加快速度吗?仅供参考,我有 2 台机器运行 13 个客户端(同一个消费者组,在一个线程中)。
    • 你对每条记录的消费者有一些耗时的操作吗?单个主题的每秒大约新消息率是多少?我试图了解您的消费者如此缓慢的原因。您可以通过将分区数量增加两倍来加速它,因此理论上您将消耗两倍。
    • 主线程消费者将数据复制到本地,其他线程更新到某个 3rd 方存储。因此,理想情况下,我们不会对从 kafka 获取数据的客户端线程进行任何耗时的操作。
    猜你喜欢
    • 2015-01-23
    • 1970-01-01
    • 2014-10-30
    • 1970-01-01
    • 2016-10-12
    • 1970-01-01
    • 1970-01-01
    • 2018-05-23
    • 2020-07-08
    相关资源
    最近更新 更多