【问题标题】:KafkaSpout very high latencyKafkaSpout 非常高的延迟
【发布时间】:2015-10-25 01:02:32
【问题描述】:

我是 apachestorm 和 kafka 的新手,作为 POC 的一部分,我正在尝试使用 Kafka 和 apachestorm 处理消息流。我正在使用来自https://github.com/apache/storm/tree/master/external/storm-kafka 的storm-kafka 源,我能够创建一个示例程序,它使用KafkaSpout 从kafka 主题读取消息并将其输出到另一个kafka 主题。我有 3 个节点 kafka(三个节点都在同一台服务器上运行)集群并创建了 8 个分区的主题。我将 KafkaSpout 并行度设置为 8,将 Bolt 的并行度设置为 8,尝试使用 8 个执行器和任务。我已经尝试在 kafka 级别、SpoutConfig 级别和 Storm 级别设置很多 tunnig 参数,但是我遇到了非常高的整体延迟问题。我需要消息处理保证,所以确实需要确认。 Storm集群有1个supervisor,zookeeper有3个noed,kafka和storm共享。它运行在具有 144MB RAM 和 16CPU 的 Red Hat Linux 机器上。使用下面的参数,我会得到非常高的 spout 进程延迟,大约 40 秒,我需要得到大约 50K 的消息/秒级别,请你帮我配置实现它。我在各个网站上浏览了很多帖子,并尝试了很多调整选项,但都没有结果。

Storm config
topology.receiver.buffer.size=16
topology.transfer.buffer.size=4096
topology.executor.receive.buffer.size=16384
topology.executor.send.buffer.size=16384
topology.spout.max.batch.size=65536
topology.max.spout.pending=10000
topology.acker.executors=20

Kafka config
fetch.size.bytes=1048576
socket.timeout.ms=10000
fetch.max.wait=10000
buffer.size.bytes=1048576

提前致谢。

风暴界面截图

【问题讨论】:

    标签: apache-kafka apache-storm


    【解决方案1】:

    您的拓扑有几个问题:

    1. 你应该有与 kafka 相同数量的 spout executors 分区
    2. 您的拓扑处理元组的速度不够快。我是 对元组如何没有因超时而开始失败感到惊讶。用一个 topology.max.spout.pending 的合理值,我推荐 150 或
      1. 这只会防止超时,您的 spout 会慢慢消耗元组,因为拓扑的其余部分无法处理它。
    3. 您需要为螺栓添加更多执行器,唯一能让您的拓扑变得更快的就是让更多执行单元发挥作用。执行器和线程不是一回事,需要在拓扑中放更多的执行器。您的单个执行器延迟为 0,097,这意味着您的单个执行器每秒可以处理大约 10309 个元组;也就是说,要达到每秒 50k 的目标,您需要至少有 5 个执行者。我确信使用您的 16 cpu 机器,您可以使用超过 1 个 CPU 来处理 Bolt。
    4. 任务的主要目的是在重新平衡期间将它们提升为执行者;因此 num tasks >= num executors。
    5. 如果您使用全局分组,则需要重新设计拓扑以使用类似于字段分组的方式。

    【讨论】:

    • 非常感谢您提供解释,它澄清了疑问。我能够通过增加执行器并将最大挂起的 spout 设置为 150 来限制 spout 速度,并且能够降低延迟。我没有使用任何全局分组,我正在使用 localorshuffle 分组。我正在做一些性能测试,以找出处理 50K msg/sec 所需的条件
    • 您的机器有 16 个内核,因此您应该可以在现场(与 kafka 分区的数量相同)和 spout 上增加至少 5 或 6 个执行器的数量(可能您需要更多不要)不敢使用 10 或 20)。
    【解决方案2】:

    查看您的 UI 屏幕截图,您的 spout 似乎发出了更多可以由您的 bolt 处理的数据。两个 spout 都发出了大约 500K 消息,但只有 250k 得到了确认(同样可以通过执行的 bolt 元组数推断出来——大约是 480K,是两个 spout 发出的元组的一半)。从一开始,40 秒的延迟是否相同?或者延迟是否会随着时间的推移而增加?如果它随着时间的推移而增加,很明显你的螺栓是瓶颈。你有两个选择:

    1. 增加螺栓和/或的平行度
    2. 设置参数spout.max.pending来限制spout输出速率

    第一个选项只有在你有足够的内核时才有意义(但到目前为止这应该不是问题,因为你提到了 16 个可用的 CPU)。 如果第二个选项适用于您,则取决于您想要实现的吞吐量。您提到了 50K msg/sec,但 UI 没有显示当前的吞吐量数(即 spout 输出速率),因此我无法判断是否可以选择节流。此外,您必须通过试错来确定spout.max.pending 的最佳值(从1000 的值开始对我来说似乎是合理的)。

    【讨论】:

    • 非常感谢您的快速回复。在storm UI中,发出和转移的数量显示非常高,但确认的数量显示很少,我验证了确认的数量是在主题上收到的消息数,我不知道为什么emmited/transferred计数显示很高。当仍然没有消息时,我看到这个数字一直在上升,这是否意味着有问题,请您指出正确的方向。
    • 关键是,spout 发出元组的速度比 bolt 处理的快。因此,元组在 Storm 内部被缓冲。这个内部队列随着时间的推移而增长。当 Kafka 没有更多可用数据时,spout 停止发出,但 bolt 继续处理,直到所有内部缓冲的元组都得到进程。内部 Storm 队列中元组的等待时间增加了测量的延迟。因此,致命的方式并没有错。您只需要限制 spouts 输出速率或增加 bolt 的 dop 即可消除这种影响。
    【解决方案3】:

    我不知道您的问题是否已解决,但除了根据您的延迟要求调整 topology.max.spout.pending 之外,您还需要调整批量大小。 将topology.spout.max.batch.size 设置为较小的数字可能有助于减少延迟。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-09-18
      • 2020-08-27
      • 2019-03-08
      • 1970-01-01
      • 2023-01-18
      • 2017-09-19
      • 2016-12-10
      • 1970-01-01
      相关资源
      最近更新 更多