【问题标题】:How to measure latency and throughput in a Storm topology如何在 Storm 拓扑中测量延迟和吞吐量
【发布时间】:2016-12-14 18:45:30
【问题描述】:

我正在通过示例 ExclamationTopology 学习 Storm。我想测量一个螺栓的延迟(将!!! 添加到一个单词所需的时间)和吞吐量(例如,每秒有多少单词通过一个螺栓)。

从here,我可以统计出一个bolt的字数和执行了多少次:

_countMetric = new CountMetric();
_wordCountMetric = new MultiCountMetric();


context.registerMetric("execute_count", _countMetric, 5);
context.registerMetric("word_count", _wordCountMetric, 60);

我知道 Storm UI 提供了 Process Latency 和 Execute Latency,而这个 post 很好地解释了它们是什么。

但是,我想记录每个螺栓的每次执行的延迟,并使用此信息与 word_count 一起计算吞吐量。

如何使用Storm Metrics 来完成此操作?

【问题讨论】:

    标签: java performance clojure cloud apache-storm


    【解决方案1】:

    虽然您的问题直截了当,并且肯定会引起很多人的兴趣,但它的答案并不像应有的那么微不足道。首先,我们需要澄清,我们真正想要测量的究竟是什么。吞吐量和延迟是术语,很容易理解,但在 Storms 分布式环境中事情变得更加复杂。

    正如出色的blog post 中所描述的,每个 Storm 主管至少有 3 个线程来完成不同的任务。当 Worker Receiver Thread 等待传入的数据元组并将它们聚合成一个块时,它们被发送到 Worker Executor Thread。这包含用户逻辑(在您的情况下为 ExclamationBolt 和一个负责传出消息的发送者。最后,在每个主管节点上,有一个 Worker Send Thread 聚合来自所有执行者,聚合它们并将它们发送到网络。

    当然,每个线程都有自己的延迟和吞吐量。对于发送者和接收者线程,它们在很大程度上取决于缓冲区大小,您可以对其进行调整。在您的情况下,您只想测量一个(执行)bolt 的延迟和吞吐量 - 这是可能的,但请记住,其他线程会对这个 bolt 产生影响。

    我的方法: 为了获得延迟和吞吐量,我使用了旧的Storm Builtin Metrics。因为我发现文档不是很清楚,所以我在这里画一条线:我们不使用新的Storm Metric API v2,我们不使用Cluster Metrics。 p>

    1. 在您的storm.yaml 中添加以下内容来激活风暴记录:
    topology.metrics.consumer.register:
      - class: "org.apache.storm.metric.LoggingMetricsConsumer"
        parallelism.hint: 1
    
    1. 您可以设置报告间隔:topology.builtin.metrics.bucket.size.secs: 10

    2. 运行您的查询。所有指标每 10 秒记录在特定的指标日志文件中。找到这个日志文件并非易事。 Storm 创建一个LoggingMetricsConsumer-Bolt 并在集群中分发它。在此节点上,您应该在 Storm 日志中找到相应的指标文件。

    3. 此指标文件包含每个执行程序的指标,您正在寻找,例如:complete-latency、execute-latency 等等。对于吞吐量,我将使用包含例如:arrival_rate_secs 的队列指标作为每秒插入多少元组的估计值。照顾在每个主管上执行的多个线程。

    祝你好运!

    【讨论】:

      猜你喜欢
      • 2022-12-25
      • 2016-11-16
      • 2017-10-29
      • 2015-04-16
      • 2018-12-14
      • 2021-04-16
      • 2019-08-27
      • 2017-11-19
      • 2022-09-23
      相关资源
      最近更新 更多