【问题标题】:Storm message failed风暴消息失败
【发布时间】:2017-06-08 15:41:36
【问题描述】:

最近我遇到了一个非常奇怪的问题。风暴集群有 3 台机器。拓扑结构是这样的,Kafka Spout A -> Bolt B -> Bolt C。我已经确认了每个bolt中的所有元组,即使内部bolt可能抛出异常(在bolt执行方法中我尝试捕获所有异常,最后确认元组)。 但是这里发生了奇怪的事情。我打印了 spout 的日志,在一台机器上打印了 spout 确认的所有元组,但在其他 2 台机器上,几乎所有元组都失败了。 60 秒后,元组一次又一次地重播。 “几乎”意味着在开始时,其他 2 台机器上的所有元组都失败了。一段时间后,两台机器上确认了少量的元组。

由于超时,元组肯定会失败。但我真的不知道他们为什么超时。根据我打印的日志,我真的确定所有元组在每个螺栓的执行方法末尾都得到了确认。所以我想知道为什么有些元组在两台机器上失败了。

我可以做些什么来找出拓扑或风暴集群出了什么问题?非常感谢并希望得到您的回复。

【问题讨论】:

    标签: apache-storm


    【解决方案1】:

    您的问题与 StormTopology 中 KafkaSpout 对背压的处理有关。

    您可以通过在拓扑配置中设置 maxSpoutPending 值来处理 KafkaSpout 的背压,

    Config config = new Config();
    config.setMaxSpoutPending(200); 
    config.setMessageTimeoutSecs(100);
    
    StormSubmitter.submitTopology("testtopology", config, builder.createTopology());
    

    maxSpoutPending 是在给定时间可以在拓扑中等待确认的元组数。设置此属性,将提示 KafkaSpout 不再使用来自 Kafka 的任何数据,除非未确认的元组计数小于 maxSpoutPending 值。

    此外,请确保您可以微调 Bolt,使其尽可能轻量级,以便元组在超时之前得到确认。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-26
      • 1970-01-01
      • 1970-01-01
      • 2019-07-07
      • 2019-01-30
      • 1970-01-01
      相关资源
      最近更新 更多